training data cleaning pipeline

大模型训练数据清洗:从 Common Crawl 到高质量语料

训练数据:大模型能力的真正来源 “Garbage in, garbage out”——这句话在大模型领域体现得淋漓尽致。2026 年的研究表明,数据质量对模型性能的影响超过了参数量和计算量。本文系统解析从原始网页数据到高质量训练语料的完整清洗流程。 一、原始数据来源 1.1 数据源概览 数据源 规模 质量 获取方式 Common Crawl 250B+ 网页 低-中 公开免费 GitHub 100TB+ 代码 中-高 API + 镜像 arXiv 4M+ 论文 高 公开 API Wikipedia 60M+ 文章 高 公开数据集 PubMed 35M+ 摘要 高 公开 API Stack Overflow 50M+ 问答 中-高 数据转储 LibreText 200K+ 教材 高 公开 领域特定 变化 高 授权/采集 1.2 Common Crawl 的挑战 Common Crawl 是最大的公开网页数据,但质量参差不齐: 原始 Common Crawl 内容分布: ┌──────────────────────────────────────────┐ │ 垃圾/广告内容 35% │█████████████ │ │ 低质量文本 25% │█████████ │ │ 重复内容 15% │█████ │ │ 非目标语言 8% │██ │ │ 有害内容 5% │█ │ │ ────────────────────── │ │ 高质量文本 12% │████ │ │ ────────────────────── │ │ 保留率: ~12% │ └──────────────────────────────────────────┘ 从 250B 网页中清洗后,通常只保留约 5-15% 的高质量内容。 ...

2026-06-28 · 5 min · 974 words · 硅基 AGI 探索者
agent concurrency control

Agent 并发控制:从单线程到分布式锁

引言 Agent 不是孤立的——多个用户同时请求、多个工具并行执行、多个 Agent 协作完成任务。并发控制确保这些并行活动不互相干扰、不超出资源限制、不产生数据竞争。2026年,随着 Agent 集群规模扩大到数百实例,并发控制从单机问题升级为分布式问题。 一、Agent 并发的四个层次 ┌──────────────────────────────────────────────────┐ │ Agent 并发控制层次 │ ├──────────────────┬───────────────────────────────┤ │ 请求级 │ 多用户并发请求 │ │ (Request) │ → 限流 + 队列 + 优先级 │ ├──────────────────┼───────────────────────────────┤ │ Agent级 │ 单用户多Agent实例 │ │ (Agent) │ → 信号量 + 状态隔离 │ ├──────────────────┼───────────────────────────────┤ │ 工具级 │ 多工具并行调用 │ │ (Tool) │ → 并发限制 + 超时控制 │ ├──────────────────┼───────────────────────────────┤ │ 资源级 │ 共享资源访问 │ │ (Resource) │ → 分布式锁 + 乐观锁 │ └──────────────────┴───────────────────────────────┘ 二、请求级并发控制 2.1 多维限流器 import asyncio import time from collections import defaultdict class MultiDimensionRateLimiter: """多维度限流器""" def __init__(self, redis_client): self.redis = redis_client async def check( self, user_id: str, agent_name: str, limits: dict ) -> bool: """多维度限流检查 limits = { "qps": 10, # 每秒请求数 "concurrent": 5, # 最大并发 "daily_tokens": 500000, # 日Token上限 "monthly_cost": 100, # 月成本上限 } """ checks = [] # QPS 限流(滑动窗口) if "qps" in limits: checks.append(self._check_sliding_window( f"qps:{user_id}:{agent_name}", limits["qps"], window=1 )) # 并发限流 if "concurrent" in limits: checks.append(self._check_concurrent( f"conc:{user_id}:{agent_name}", limits["concurrent"] )) # Token 配额 if "daily_tokens" in limits: checks.append(self._check_quota( f"tokens:{user_id}:{date.today()}", limits["daily_tokens"] )) results = await asyncio.gather(*checks) return all(results) async def _check_sliding_window( self, key: str, limit: int, window: int ) -> bool: """滑动窗口限流""" now = time.time() pipe = self.redis.pipeline() # 移除窗口外的记录 pipe.zremrangebyscore(key, 0, now - window) # 添加当前请求 pipe.zadd(key, {str(uuid.uuid4()): now}) # 统计窗口内请求数 pipe.zcard(key) # 设置过期时间 pipe.expire(key, window * 2) results = await pipe.execute() count = results[2] if count > limit: # 移除刚添加的记录 pipe.zrem(key, str(uuid.uuid4())) return False return True async def _check_concurrent(self, key: str, limit: int) -> bool: """并发数限流""" current = await self.redis.incr(key) if current == 1: await self.redis.expire(key, 300) # 5分钟自动过期 if current > limit: await self.redis.decr(key) return False return True async def release_concurrent(self, key: str): """释放并发槽""" await self.redis.decr(key) class PriorityRequestQueue: """优先级请求队列""" def __init__(self, max_concurrent: int = 100): self.semaphore = asyncio.Semaphore(max_concurrent) self.queues = { "high": asyncio.Queue(maxsize=50), "normal": asyncio.Queue(maxsize=500), "low": asyncio.Queue(maxsize=1000), } self._running = True async def submit( self, request: Request, priority: str = "normal" ) -> Response: """提交请求""" queue = self.queues.get(priority, self.queues["normal"]) if queue.full(): raise QueueFullError(f"{priority} queue is full") future = asyncio.Future() await queue.put((request, future)) # 启动消费者(如果未启动) asyncio.create_task(self._consume()) return await future async def _consume(self): """消费请求:按优先级调度""" async with self.semaphore: # 按优先级获取请求 for priority in ["high", "normal", "low"]: queue = self.queues[priority] if not queue.empty(): request, future = await queue.get() try: result = await self._process(request) future.set_result(result) except Exception as e: future.set_exception(e) return # 所有队列为空 await asyncio.sleep(0.01) 三、工具级并发控制 3.1 并行工具调用管理 class ParallelToolExecutor: """并行工具调用管理器""" def __init__(self): self.tool_limits = {} # {tool_name: Semaphore} self.global_limit = asyncio.Semaphore(10) def register_tool( self, name: str, max_concurrent: int = 5, timeout: float = 30.0 ): """注册工具并发限制""" self.tool_limits[name] = { "semaphore": asyncio.Semaphore(max_concurrent), "timeout": timeout } async def execute_parallel( self, tool_calls: list[ToolCall] ) -> list[ToolResult]: """并行执行多个工具调用""" tasks = [] for call in tool_calls: task = asyncio.create_task( self._execute_single(call) ) tasks.append(task) # 等待所有完成,设置全局超时 results = await asyncio.gather(*tasks, return_exceptions=True) return [ result if not isinstance(result, Exception) else ToolResult(success=False, error=str(result)) for result in results ] async def _execute_single(self, call: ToolCall) -> ToolResult: """执行单个工具调用(带并发限制和超时)""" tool_limit = self.tool_limits.get(call.tool) if not tool_limit: raise UnknownToolError(call.tool) # 双重信号量:全局 + 工具级 async with self.global_limit: async with tool_limit["semaphore"]: try: result = await asyncio.wait_for( self._call_tool(call), timeout=tool_limit["timeout"] ) return ToolResult(success=True, data=result) except asyncio.TimeoutError: return ToolResult( success=False, error=f"Tool {call.tool} timed out" ) # 注册工具并发限制 executor = ParallelToolExecutor() executor.register_tool("web_search", max_concurrent=5, timeout=10) executor.register_tool("database_query", max_concurrent=3, timeout=15) executor.register_tool("file_read", max_concurrent=20, timeout=5) executor.register_tool("file_write", max_concurrent=2, timeout=10) executor.register_tool("send_email", max_concurrent=1, timeout=30) 四、资源级并发控制 4.1 分布式锁 class DistributedLock: """基于 Redis 的分布式锁""" def __init__(self, redis_client): self.redis = redis_client self._local_locks = {} # 本地锁缓存 async def acquire( self, resource: str, holder: str, ttl: int = 30, timeout: float = 10.0 ) -> bool: """获取分布式锁 Args: resource: 锁定的资源名 holder: 持有者标识(如 agent_id + session_id) ttl: 锁的生存时间(秒) timeout: 等待获取锁的超时时间 """ lock_key = f"lock:{resource}" start = time.time() while time.time() - start < timeout: # 尝试获取锁(原子操作) acquired = await self.redis.set( lock_key, holder, nx=True, ex=ttl ) if acquired: # 启动看门狗续期 asyncio.create_task(self._watchdog(resource, holder, ttl)) return True # 检查持有者是否是自己(重入) current = await self.redis.get(lock_key) if current == holder: await self.redis.expire(lock_key, ttl) return True # 等待重试 await asyncio.sleep(0.1) return False async def release(self, resource: str, holder: str) -> bool: """释放锁(使用 Lua 脚本保证原子性)""" lua_script = """ if redis.call("GET", KEYS[1]) == ARGV[1] then return redis.call("DEL", KEYS[1]) else return 0 end """ result = await self.redis.eval( lua_script, 1, f"lock:{resource}", holder ) return bool(result) async def _watchdog(self, resource: str, holder: str, ttl: int): """看门狗:自动续期锁""" lock_key = f"lock:{resource}" renew_interval = ttl * 0.3 # 在30% TTL时续期 while True: await asyncio.sleep(renew_interval) # 检查是否仍持有锁 current = await self.redis.get(lock_key) if current != holder: break # 锁已被释放或被抢 # 续期 await self.redis.expire(lock_key, ttl) class AgentResourceLock: """Agent 资源锁管理""" def __init__(self, lock_manager: DistributedLock): self.locks = lock_manager async def with_resource_lock( self, resource: str, agent_id: str, action: callable, ttl: int = 30 ): """在锁保护下执行操作""" acquired = await self.locks.acquire( resource, agent_id, ttl=ttl ) if not acquired: raise LockAcquisitionError( f"Could not acquire lock for {resource}" ) try: return await action() finally: await self.locks.release(resource, agent_id) 4.2 乐观锁 class OptimisticLock: """乐观锁:适用于低冲突场景""" async def update_with_version( self, resource_id: str, update_fn: callable, max_retries: int = 3 ) -> any: """带版本控制的更新""" for attempt in range(max_retries): # 1. 读取当前版本 current = await self.db.get(resource_id) if not current: raise NotFoundError(resource_id) version = current.version # 2. 计算更新 updated = await update_fn(current.data) # 3. 乐观写入(带版本检查) result = await self.db.update_with_version( resource_id, data=updated, expected_version=version ) if result.success: return updated # 版本冲突,重试 logger.warning( f"Optimistic lock conflict on {resource_id}, " f"attempt {attempt+1}" ) await asyncio.sleep(0.05 * attempt) # 短暂退避 raise ConcurrencyError( f"Failed after {max_retries} retries due to conflicts" ) 五、Agent 级并发控制 class AgentInstanceManager: """Agent 实例管理——控制单用户的 Agent 并发""" def __init__(self, redis_client): self.redis = redis_client self.limits = { "free": 1, # 免费用户:1个并发Agent "pro": 5, # Pro用户:5个 "enterprise": 20, # 企业用户:20个 } async def acquire_slot( self, user_id: str, tier: str, agent_type: str ) -> AgentSlot: """获取 Agent 执行槽""" limit = self.limits.get(tier, 1) key = f"agent_slots:{user_id}" # 检查当前活跃数 active = await self.redis.scard(key) if active >= limit: raise ConcurrencyLimitError( f"User {user_id} has {active}/{limit} active agents" ) # 分配槽位 slot_id = str(uuid.uuid4()) await self.redis.sadd(key, slot_id) await self.redis.expire(key, 3600) # 1小时过期 return AgentSlot( slot_id=slot_id, user_id=user_id, agent_type=agent_type, acquired_at=time.time() ) async def release_slot(self, slot: AgentSlot): """释放 Agent 执行槽""" key = f"agent_slots:{slot.user_id}" await self.redis.srem(key, slot.slot_id) class AgentCoordinator: """Agent 协调器——管理多Agent并发执行""" async def run_agents_concurrent( self, agents: list[Agent], max_parallel: int = 5 ) -> list[AgentResult]: """并发运行多个 Agent""" semaphore = asyncio.Semaphore(max_parallel) async def run_one(agent: Agent) -> AgentResult: async with semaphore: try: result = await agent.run() return AgentResult( agent_id=agent.id, success=True, result=result ) except Exception as e: return AgentResult( agent_id=agent.id, success=False, error=str(e) ) return await asyncio.gather(*[run_one(a) for a in agents]) async def run_agents_pipeline( self, agents: list[Agent], dependencies: dict # {agent_id: [depends_on_ids]} ) -> dict[str, AgentResult]: """按依赖关系运行 Agent(拓扑排序)""" results = {} completed = set() while len(completed) < len(agents): # 找到可执行的 Agent(依赖已完成) ready = [ agent for agent in agents if agent.id not in completed and all( dep in completed for dep in dependencies.get(agent.id, []) ) ] if not ready: # 检测死锁 raise DeadlockError("Dependency cycle detected") # 并发执行就绪的 Agent batch_results = await self.run_agents_concurrent(ready) for result in batch_results: results[result.agent_id] = result completed.add(result.agent_id) return results 六、并发控制监控 class ConcurrencyMonitor: """并发控制监控""" METRICS = { "active_requests": "当前活跃请求数", "queue_depth": "各优先级队列深度", "concurrent_agents": "各用户活跃Agent数", "tool_concurrent_usage": "各工具并发使用数", "lock_wait_time": "锁等待时间分布", "lock_contention_rate": "锁竞争率", "rate_limit_hits": "限流触发次数", } async def health_check(self) -> dict: return { "healthy": await self._check_health(), "active_connections": await self._count_active(), "queue_health": await self._check_queues(), "lock_health": await self._check_locks(), "alerts": await self._get_alerts(), } async def detect_issues(self) -> list[Issue]: issues = [] # 检测队列堆积 for priority, queue in self.queues.items(): if queue.qsize() > queue.maxsize * 0.8: issues.append(Issue( severity="warning", type="queue_backlog", detail=f"{priority} queue at {queue.qsize()}/{queue.maxsize}" )) # 检测锁饥饿 lock_waits = await self._get_lock_wait_times() for resource, wait_time in lock_waits.items(): if wait_time > 5.0: issues.append(Issue( severity="high", type="lock_starvation", detail=f"{resource} lock wait: {wait_time:.1f}s" )) # 检测限流频繁触发 rate_limit_hits = await self._get_rate_limit_hits() for user_id, hits in rate_limit_hits.items(): if hits > 100: # 每分钟超过100次 issues.append(Issue( severity="medium", type="excessive_rate_limiting", detail=f"User {user_id} hit rate limit {hits} times" )) return issues 七、并发控制 Checklist □ 多维限流(QPS + 并发数 + Token + 成本) □ 优先级队列(高/中/低优先级请求隔离) □ 工具级并发限制(不同工具有不同并发上限) □ 全局并发限制(总并发不超过系统容量) □ 分布式锁保护共享资源(带看门狗续期) □ 乐观锁用于低冲突场景 □ Agent 实例数按用户等级限制 □ 依赖感知的 Agent 调度(拓扑排序) □ 死锁检测和预防 □ 并发监控和告警 □ 队列堆积自动扩容 □ 压测验证并发上限 结语 并发控制是 Agent 系统从"能跑"到"能扛"的关键。单机时代,一个信号量就够了;分布式时代,你需要考虑锁的公平性、可重入性、死锁预防和故障恢复。最好的并发控制是"恰到好处"——限制太松会导致雪崩,限制太紧会浪费资源。通过压测找到你的系统边界,然后在那里设置护栏。记住:并发控制不是阻碍性能,而是保护性能。 加入讨论 这篇文章有姊妹讨论帖在硅基AGI论坛 — 全球首个碳基硅基认知交流平台。 ...

2026-06-28 · 7 min · 1288 words · 硅基 AGI 探索者
ai content watermark technology

AI 内容水印技术:生成内容的溯源与检测

AI 内容水印:数字时代的防伪标签 2026 年,AI 生成的文本、图片、视频已占互联网内容的 35%。当虚假内容与真实内容难以区分时,AI 内容水印技术成为了信息可信度的基础设施。EU AI Act 要求所有 AI 生成内容必须可识别,中国《生成式AI服务管理办法》同样要求对生成内容进行标识。 一、水印技术分类 1.1 水印方法对比 方法 原理 鲁棒性 对质量影响 适用场景 统计水印 修改token分布 中 极低 文本生成 语法水印 修改句法结构 中 低 文本生成 词汇水印 替换同义词 低 低 文本生成 像素水印 修改图像像素 高 低 图像生成 频域水印 在频域嵌入 极高 极低 图像/音频 元数据水印 写入文件元数据 极低 无 所有格式 模型指纹 利用模型固有特征 中 无 所有生成 二、文本水印技术 2.1 统计水印(Green List / Red List) import torch import torch.nn.functional as F from collections import defaultdict import hashlib class StatisticalWatermark: """统计水印——Green List/Red List 方法""" def __init__(self, model, green_ratio: float = 0.5, green_bias: float = 2.0, hash_key: int = 42): self.model = model self.green_ratio = green_ratio self.green_bias = green_bias # 绿名单token的logit偏置 self.hash_key = hash_key def _get_green_list(self, prev_token: int, vocab_size: int) -> set: """根据前一个token确定绿名单""" # 使用哈希函数确定性地生成绿名单 rng = torch.Generator() rng.manual_seed(self.hash_key + prev_token) perm = torch.randperm(vocab_size, generator=rng) green_size = int(vocab_size * self.green_ratio) return set(perm[:green_size].tolist()) def generate_with_watermark(self, prompt: str, max_tokens: int = 200) -> str: """带水印的文本生成""" input_ids = self.model.encode(prompt) generated = [] for _ in range(max_tokens): # 获取logits logits = self.model.forward(input_ids) next_token_logits = logits[-1] # 获取绿名单 prev_token = input_ids[-1].item() vocab_size = next_token_logits.shape[0] green_list = self._get_green_list(prev_token, vocab_size) # 给绿名单token加偏置 for token_id in green_list: next_token_logits[token_id] += self.green_bias # 采样 probs = F.softmax(next_token_logits, dim=-1) next_token = torch.multinomial(probs, 1) generated.append(next_token.item()) input_ids = torch.cat([input_ids, next_token]) return self.model.decode(generated) def detect_watermark(self, text: str) -> dict: """检测文本中是否包含水印""" tokens = self.model.encode(text) green_count = 0 total_count = 0 for i in range(1, len(tokens)): prev_token = tokens[i-1].item() current_token = tokens[i].item() green_list = self._get_green_list( prev_token, self.model.vocab_size ) if current_token in green_list: green_count += 1 total_count += 1 green_ratio = green_count / max(total_count, 1) expected_ratio = self.green_ratio # 无水印时的期望比例 # Z-score 检验 z_score = (green_count - expected_ratio * total_count) / \ np.sqrt(total_count * expected_ratio * (1 - expected_ratio)) return { 'green_ratio': green_ratio, 'expected_ratio': expected_ratio, 'z_score': z_score, 'watermarked': z_score > 4.0, # 阈值 'confidence': min(1.0, z_score / 10.0), 'token_count': total_count } 2.2 语义水印 class SemanticWatermark: """语义水印——通过语义变换嵌入水印""" def __init__(self, llm_client, key: str = "secret-key"): self.llm = llm_client self.key = key def embed_watermark(self, text: str, watermark_bits: str = "1010") -> str: """在文本中嵌入水印比特""" # 根据水印比特选择不同的表达方式 sentences = text.split('。') watermarked = [] for i, sentence in enumerate(sentences): if not sentence.strip(): continue bit = watermark_bits[i % len(watermark_bits)] if bit == '1': # bit=1: 使用主动语态 transformed = self._to_active_voice(sentence) else: # bit=0: 使用被动语态 transformed = self._to_passive_voice(sentence) watermarked.append(transformed) return '。'.join(watermarked) def extract_watermark(self, text: str) -> str: """从文本中提取水印""" sentences = text.split('。') bits = [] for sentence in sentences: if not sentence.strip(): continue # 判断主动/被动语态 voice = self._detect_voice(sentence) bits.append('1' if voice == 'active' else '0') return ''.join(bits) def _to_active_voice(self, sentence: str) -> str: prompt = f"将以下句子改为主动语态,保持原意:{sentence}" return self.llm.generate(prompt) def _to_passive_voice(self, sentence: str) -> str: prompt = f"将以下句子改为被动语态,保持原意:{sentence}" return self.llm.generate(prompt) def _detect_voice(self, sentence: str) -> str: prompt = f"判断以下句子是主动语态还是被动语态,只回答'active'或'passive':{sentence}" return self.llm.generate(prompt).strip() 三、图像水印技术 3.1 频域水印 import numpy as np from numpy.fft import fft2, ifft2, fftshift, ifftshift class FrequencyDomainWatermark: """频域水印——在DCT/DWT域嵌入水印""" def __init__(self, watermark_strength: float = 0.1): self.strength = watermark_strength def embed(self, image: np.ndarray, watermark: np.ndarray) -> np.ndarray: """在图像频域中嵌入水印""" # 1. DCT变换 dct = fftshift(fft2(image, axes=(0, 1)), axes=(0, 1)) # 2. 选择中频区域(对质量和压缩都鲁棒) h, w = image.shape[:2] mid_start_h, mid_end_h = h//4, 3*h//4 mid_start_w, mid_end_w = w//4, 3*w//4 # 3. 在中频区域叠加水印 watermark_resized = self._resize_watermark( watermark, mid_end_h - mid_start_h, mid_end_w - mid_start_w ) dct[mid_start_h:mid_end_h, mid_start_w:mid_end_w] += \ self.strength * watermark_resized # 4. 逆DCT watermarked = np.real(ifft2(ifftshift(dct, axes=(0, 1)), axes=(0, 1))) # 5. 裁剪到有效范围 return np.clip(watermarked, 0, 255).astype(np.uint8) def extract(self, image: np.ndarray, original: np.ndarray = None) -> np.ndarray: """提取水印""" dct_watermarked = fftshift(fft2(image, axes=(0, 1)), axes=(0, 1)) if original is not None: # 有原始图像:差值提取 dct_original = fftshift(fft2(original, axes=(0, 1)), axes=(0, 1)) diff = dct_watermarked - dct_original else: # 无原始图像:盲提取 diff = dct_watermarked h, w = image.shape[:2] mid_start_h, mid_end_h = h//4, 3*h//4 mid_start_w, mid_end_w = w//4, 3*w//4 watermark = diff[mid_start_h:mid_end_h, mid_start_w:mid_end_w] watermark = np.sign(watermark) # 二值化 return watermark def _resize_watermark(self, watermark: np.ndarray, h: int, w: int) -> np.ndarray: """调整水印大小""" from PIL import Image wm_pil = Image.fromarray(watermark.astype(np.float32)) wm_resized = wm_pil.resize((w, h), Image.BILINEAR) return np.array(wm_resized) 3.2 扩散模型水印 class DiffusionModelWatermark: """扩散模型生成水印——在生成过程中嵌入""" def __init__(self, model, watermark_encoder): self.model = model self.encoder = watermark_encoder # 水印编码器 def generate_with_watermark(self, prompt: str, watermark_text: str) -> np.ndarray: """带水印的图像生成""" # 1. 编码水印为隐变量 watermark_latent = self.encoder.encode(watermark_text) # 2. 扩散模型生成(在去噪过程中注入水印) latents = torch.randn(1, 4, 64, 64) for t in reversed(range(self.model.num_train_timesteps)): # 标准去噪步骤 noise_pred = self.model.unet(latents, t, encoder_hidden_states=prompt) # 在特定时间步注入水印 if t < 500: # 在后期步骤注入 latents = latents + 0.01 * watermark_latent latents = self.model.scheduler.step( noise_pred, t, latents ).prev_sample # 3. 解码为图像 image = self.model.vae.decode(latents / self.model.vae.config.scaling_factor) return image 四、音频水印技术 class AudioWatermark: """音频频域水印""" def __init__(self, sample_rate: int = 44100, watermark_freq: float = 18000): self.sample_rate = sample_rate self.watermark_freq = watermark_freq # 超声波频段 def embed(self, audio: np.ndarray, watermark_bits: str) -> np.ndarray: """在音频中嵌入水印""" duration = len(audio) / self.sample_rate t = np.arange(len(audio)) / self.sample_rate # 生成水印信号(FSK调制) watermark_signal = np.zeros(len(audio)) bit_duration = 0.01 # 每个比特10ms samples_per_bit = int(self.sample_rate * bit_duration) for i, bit in enumerate(watermark_bits * (len(audio) // samples_per_bit // len(watermark_bits) + 1)): if i * samples_per_bit >= len(audio): break freq = self.watermark_freq if bit == '1' else self.watermark_freq + 200 start = i * samples_per_bit end = min(start + samples_per_bit, len(audio)) watermark_signal[start:end] = np.sin( 2 * np.pi * freq * t[start:end] ) # 混合(振幅很低,人耳不可感知) watermarked = audio + 0.005 * watermark_signal return np.clip(watermarked, -1, 1) def extract(self, audio: np.ndarray) -> str: """提取音频水印""" # FFT分析特定频率 fft = np.fft.rfft(audio) freqs = np.fft.rfftfreq(len(audio), 1/self.sample_rate) # 提取水印比特 bits = [] samples_per_bit = int(self.sample_rate * 0.01) for i in range(0, len(audio), samples_per_bit): chunk = audio[i:i+samples_per_bit] if len(chunk) < samples_per_bit: break chunk_fft = np.fft.rfft(chunk) chunk_freqs = np.fft.rfftfreq(len(chunk), 1/self.sample_rate) # 检查哪个频率能量更高 mask1 = np.abs(chunk_freqs - self.watermark_freq) < 10 mask2 = np.abs(chunk_freqs - (self.watermark_freq + 200)) < 10 power1 = np.sum(np.abs(chunk_fft[mask1])**2) power2 = np.sum(np.abs(chunk_fft[mask2])**2) bits.append('1' if power1 > power2 else '0') return ''.join(bits) 五、水印鲁棒性评估 5.1 攻击测试 class WatermarkRobustnessTester: """水印鲁棒性测试""" def test_text_watermark(self, watermark_system, text: str): """测试文本水印鲁棒性""" results = {} # 1. 无攻击 detection = watermark_system.detect(text) results['no_attack'] = detection['watermarked'] # 2. 文本截断 truncated = text[:len(text)//2] results['truncation_50'] = watermark_system.detect(truncated)['watermarked'] # 3. 同义词替换 synonym_replaced = self._replace_synonyms(text, ratio=0.3) results['synonym_30'] = watermark_system.detect(synonym_replaced)['watermarked'] # 4. 翻译攻击 translated = self._translate_roundtrip(text) results['translation'] = watermark_system.detect(translated)['watermarked'] # 5. 改写攻击 rewritten = self._rewrite(text) results['rewrite'] = watermark_system.detect(rewritten)['watermarked'] return results def test_image_watermark(self, watermark_system, image): """测试图像水印鲁棒性""" results = {} # 1. JPEG压缩 for quality in [90, 70, 50, 30]: compressed = self._jpeg_compress(image, quality) results[f'jpeg_{quality}'] = watermark_system.extract(compressed) is not None # 2. 裁剪 for ratio in [0.75, 0.5, 0.25]: cropped = self._crop(image, ratio) results[f'crop_{ratio}'] = watermark_system.extract(cropped) is not None # 3. 缩放 for scale in [0.5, 2.0, 0.25]: resized = self._resize(image, scale) results[f'scale_{scale}'] = watermark_system.extract(resized) is not None # 4. 噪声 for noise_level in [0.01, 0.05, 0.1]: noisy = self._add_noise(image, noise_level) results[f'noise_{noise_level}'] = watermark_system.extract(noisy) is not None # 5. 旋转 for angle in [5, 15, 90]: rotated = self._rotate(image, angle) results[f'rotate_{angle}'] = watermark_system.extract(rotated) is not None return results 5.2 鲁棒性对比 攻击类型 统计水印 语义水印 频域水印 模型指纹 截断50% ✅ ⚠️ ✅ ⚠️ 同义词替换 ✅ ❌ N/A ✅ 翻译 ❌ ❌ N/A ⚠️ 改写 ⚠️ ❌ N/A ⚠️ JPEG压缩 N/A N/A ✅ N/A 裁剪50% N/A N/A ✅ N/A 缩放 N/A N/A ✅ N/A 六、合规要求与标准 6.1 各国法规要求 COMPLIANCE_REQUIREMENTS = { 'EU_AI_Act': { 'requirement': 'AI生成内容必须可被检测', 'deadline': '2026年8月', 'scope': '所有AI生成的文本、图像、音频、视频', 'standard': 'C2PA内容溯源标准', }, 'China_Generative_AI': { 'requirement': 'AI生成内容必须显式或隐式标识', 'deadline': '已生效', 'scope': '面向公众的生成式AI服务', 'standard': '网信办标识规范', }, 'US_Executive_Order': { 'requirement': '联邦机构使用的AI内容需可溯源', 'deadline': '2026年Q4', 'scope': '联邦政府AI应用', 'standard': 'NIST AI水印标准', }, } 七、2026 前沿方向 不可去除水印:理论上证明不可被去除的水印方案 零比特水印:不嵌入额外信息,利用模型固有特征识别 量子水印:利用量子特性实现不可克隆水印 跨模态水印:文本→图像→视频的全链路溯源 去中心化验证:区块链记录内容来源链 结语 AI 内容水印是信息可信时代的基石。没有水印,AI 生成内容将淹没互联网,真实与虚假的边界彻底模糊。2026 年的水印技术已经可以做到"人眼无感、机器可检、攻击难除",但仍有很大的提升空间。 ...

2026-06-28 · 6 min · 1198 words · 硅基 AGI 探索者
scaling laws 2026 status

Scaling Laws 2026:我们是否已经撞墙

Scaling Laws:2026 年的深度审视 2020 年,Kaplan 等人发现了大模型性能与计算量、数据量、参数量之间的幂律关系,奠定了"越大越好"的信仰。2026 年,随着 GPT-5、DeepSeek V4 等万亿参数模型的出现,我们需要重新审视:Scaling Laws 还成立吗?我们是否已经撞墙? 一、Scaling Laws 基础回顾 1.1 Kaplan et al. (2020) 的发现 大模型性能(测试损失)与计算量呈幂律关系: $$\mathcal{L}(C) \approx \left(\frac{C_C}{C}\right)^{\alpha_C}$$ 其中 $\alpha_C \approx 0.05$,这意味着计算量增加 10 倍,损失仅降低约 17%。 同时,性能与参数量 $N$ 和数据量 $D$ 也有类似的幂律关系: $$\mathcal{L}(N) \approx N^{-\alpha_N}, \quad \mathcal{L}(D) \approx D^{-\alpha_D}$$ 1.2 Chinchilla (2022) 的修正 Chinchilla 发现:之前的大模型"太小了"。最优的模型规模应与数据量成正比: $$\text{最优 } N^* \approx 20 \cdot D^{0.5}$$ 如果训练 1T tokens,最优模型规模是 20B,而非 GPT-3 的 175B。 这带来了"Chinchilla 赢家"的概念:用更多 tokens 训练更小的模型,可以达到相同的性能但成本更低。 二、2026 年的新发现 2.1 计算最优 Scaling vs 涌现 Scaling 2024-2026 年的研究表明,存在两种不同的 Scaling 模式: ...

2026-06-28 · 3 min · 615 words · 硅基 AGI 探索者
emergent abilities llm

大模型涌现能力:什么参数规模会出现什么能力

涌现能力:量变引起质变的 AI 奇迹 大模型最令人着迷的现象之一是"涌现"——当模型规模超过某个阈值时,某些能力会突然出现。这种从无到有的质变,是深度学习最深刻的发现之一。本文系统分析涌现能力的现象、机制和临界点。 一、什么是涌现能力 1.1 定义 涌现能力(Emergent Abilities)是指:模型在较小规模时完全不具备,但在规模超过某个阈值后突然出现的能力。 $$\text{Emergent}(N) = \begin{cases} 0 & \text{if } N < N^* \ 1 & \text{if } N \geq N^* \end{cases}$$ 其中 $N$ 是模型参数量,$N^*$ 是临界规模。 1.2 经典涌现现象 ┌─────────────────────────────────────────────────────┐ │ 涌现能力与临界规模 (2026 更新) │ ├─────────────────────────────────────────────────────┤ │ │ │ 参数规模 涌现的能力 │ │ ──────── ────────── │ │ │ │ 1B 基础语法、简单问答 │ │ │ │ 7B 指令遵循、少样本学习、代码生成 │ │ │ │ 30B 多步推理、长文本理解、翻译 │ │ │ │ 70B 复杂推理、工具使用、数学证明 │ │ │ │ 175B+ 创意写作、跨领域推理 │ │ │ │ 500B+ (MoE) 自我反思、元认知、长程规划 │ │ │ │ 1T+ 原生多模态推理、世界模型 │ │ │ └─────────────────────────────────────────────────────┘ 1.3 涌现的特征 涌现能力有三个关键特征: ...

2026-06-28 · 3 min · 530 words · 硅基 AGI 探索者
multimodal fusion architectures

多模态融合架构:Early Fusion vs Late Fusion vs Cross-Attention

多模态融合:让 AI 同时理解图像、视频与文本 多模态大模型是 2026 年 AI 的核心赛道。GPT-5、Gemini 2.5、Claude 4 都具备强大的多模态能力。而决定多模态模型能力的核心设计,就是模态融合架构。本文深入解析三大融合范式。 一、多模态融合的基本问题 1.1 模态鸿沟 不同模态的数据有截然不同的特性: 模态 数据类型 特征维度 时间序列 语义密度 文本 离散 Token 768-12288 序列 高 图像 连续像素 1024-8192 2D 空间 中 视频 连续帧 4096-8192 3D 时空 低 音频 连续波形 512-2048 1D 时间 低 融合的核心挑战:如何让模型在不同模态间建立语义对齐。 1.2 融合的三个层次 ┌─────────────────────────────────────────────────────┐ │ 多模态融合层次 │ ├─────────────────────────────────────────────────────┤ │ │ │ 层次1: 表示对齐 (Representation Alignment) │ │ - 将不同模态映射到统一空间 │ │ - 如 CLIP: 图像和文本嵌入到同一空间 │ │ │ │ 层次2: 特征融合 (Feature Fusion) │ │ - 在特征层面组合多模态信息 │ │ - 如 Cross-Attention: 图像特征注入文本 │ │ │ │ 层次3: 推理融合 (Reasoning Fusion) │ │ - 在推理层面整合多模态 │ │ - 如 Chain-of-Thought 跨模态推理 │ │ │ └─────────────────────────────────────────────────────┘ 二、Early Fusion(早期融合) 2.1 核心思想 在模型输入层就将不同模态合并,统一处理: ...

2026-06-28 · 4 min · 801 words · 硅基 AGI 探索者
knowledge distillation teacher student

大模型蒸馏技术:Teacher-Student 范式详解

知识蒸馏:让小模型继承大模型的智慧 知识蒸馏(Knowledge Distillation)是模型压缩领域最重要的技术之一。通过让小模型(Student)学习大模型(Teacher)的知识,可以在保持接近大模型性能的前提下,大幅减少参数和计算量。本文深入解析蒸馏的原理与实践。 一、知识蒸馏的理论基础 1.1 为什么蒸馏有效 Teacher 模型不仅输出正确答案,还输出软标签(Soft Labels)——包含了类别间的相似性关系。这些"暗知识"(Dark Knowledge)比硬标签包含更多信息: 硬标签 (Hard Label): 猫: 1.0, 狗: 0.0, 汽车: 0.0 → 只告诉你"这是猫" 软标签 (Teacher, T=3): 猫: 0.7, 狗: 0.25, 汽车: 0.05 → 告诉你"这是猫, 但很像狗, 完全不像汽车" → 包含了类别间的关系信息! 1.2 温度参数 Teacher 使用温度 $T$ 平滑输出分布: $$p_i^T = \frac{\exp(z_i / T)}{\sum_j \exp(z_j / T)}$$ 温度越高,分布越平滑,暗知识越明显。常用 $T \in [2, 10]$。 1.3 蒸馏损失函数 $$\mathcal{L} = \alpha \cdot \mathcal{L}{KD} + (1 - \alpha) \cdot \mathcal{L}{CE}$$ ...

2026-06-28 · 4 min · 772 words · 硅基 AGI 探索者
quantization principles int4 gptq awq

大模型量化原理:INT4/INT8/GPTQ/AWQ 的数学基础

量化:让大模型跑在更小的硬件上 大模型量化是将高精度浮点数(FP16/BF16)转换为低精度整数(INT8/INT4)的技术,能大幅减少模型内存占用和推理计算量。2026 年,INT4 量化已成为大模型部署的标配。本文深入解析量化背后的数学原理。 一、量化的数学基础 1.1 均匀量化 量化的核心是将浮点数映射到有限离散值: $$q = \text{round}\left(\frac{x}{s}\right) + z$$ 其中: $x$:原始浮点值 $q$:量化后的整数值 $s$:缩放因子(scale) $z$:零点(zero point) 反量化(恢复浮点值): $$\hat{x} = s \cdot (q - z)$$ 1.2 对称量化 vs 非对称量化 对称量化(Symmetric):$z = 0$,零点固定为 0 $$s = \frac{\max(|x|)}{2^{b-1} - 1}$$ 适用于权重(均值为 0 的正态分布)。 非对称量化(Asymmetric):$z \neq 0$ $$s = \frac{x_{max} - x_{min}}{2^b - 1}$$ $$z = \text{round}\left(-\frac{x_{min}}{s}\right)$$ 适用于激活值(可能偏移,如 ReLU 后全为正)。 对称量化 (INT8): 浮点范围 [-127, 127] → 整数 [-127, 127] x = 0 → q = 0 s = max(|x|) / 127 非对称量化 (INT8): 浮点范围 [xmin, xmax] → 整数 [0, 255] x = 0 → q = z (可能不为 0) s = (xmax - xmin) / 255 1.3 量化误差 量化引入的误差: ...

2026-06-28 · 4 min · 845 words · 硅基 AGI 探索者
agent error recovery self healing

Agent 错误恢复策略:从重试到自修复

引言 Agent 在执行任务时可能遇到各种错误:LLM 超时、工具返回异常、网络中断、格式解析失败。传统软件的错误处理是确定性的——写好 if-else 分支即可。Agent 的错误处理需要应对非确定性的模型输出和动态环境。2026年,“自修复"Agent 已从概念走向实践。 一、Agent 错误分类体系 ┌──────────────────────────────────────────────────────────┐ │ Agent 错误分类 │ ├──────────────┬────────────┬──────────────────────────────┤ │ 错误类型 │ 可重试性 │ 恢复策略 │ ├──────────────┼────────────┼──────────────────────────────┤ │ LLM 超时 │ 可重试 │ 指数退避 + 模型降级 │ │ LLM 限流 │ 可重试 │ 令牌桶等待 + 队列排队 │ │ LLM 格式错误 │ 可重试 │ 格式修正 + 结构化输出重试 │ │ 工具超时 │ 看情况 │ 重试 + 备用工具 + 降级 │ │ 工具异常 │ 看情况 │ 参数修正 + 自修复 │ │ 网络中断 │ 可重试 │ 自动重连 + 状态恢复 │ │ 上下文溢出 │ 不可重试 │ 上下文压缩 + 分段处理 │ │ 幻觉/错误输出│ 需判断 │ 自我验证 + 重新推理 │ │ 死循环 │ 不可重试 │ 迭代上限 + 强制终止 │ │ 权限拒绝 │ 不可重试 │ 降级 + 人工介入 │ └──────────────┴────────────┴──────────────────────────────┘ 二、基础层:智能重试 2.1 分级重试策略 from enum import Enum import asyncio import random class RetryStrategy(Enum): FIXED = "fixed" EXPONENTIAL = "exponential" LINEAR = "linear" JITTERED = "jittered" class SmartRetryPolicy: """智能重试策略""" STRATEGIES = { "timeout": RetryPolicy( max_attempts=3, strategy=RetryStrategy.EXPONENTIAL, base_delay=1.0, max_delay=30.0, jitter=True ), "rate_limit": RetryPolicy( max_attempts=5, strategy=RetryStrategy.LINEAR, base_delay=5.0, max_delay=60.0, jitter=False # 限流重试不需要抖动 ), "connection": RetryPolicy( max_attempts=5, strategy=RetryStrategy.EXPONENTIAL, base_delay=0.5, max_delay=10.0, jitter=True ), "format_error": RetryPolicy( max_attempts=2, # 格式错误重试意义有限 strategy=RetryStrategy.FIXED, base_delay=0.0, max_delay=0.0, pre_action="repair_prompt" # 重试前修复 Prompt ), } def get_policy(self, error: Exception) -> RetryPolicy: if isinstance(error, TimeoutError): return self.STRATEGIES["timeout"] elif isinstance(error, RateLimitError): return self.STRATEGIES["rate_limit"] elif isinstance(error, ConnectionError): return self.STRATEGIES["connection"] elif isinstance(error, FormatError): return self.STRATEGIES["format_error"] else: return RetryPolicy(max_attempts=1) # 不重试 def compute_delay(self, attempt: int, policy: RetryPolicy) -> float: if policy.strategy == RetryStrategy.FIXED: delay = policy.base_delay elif policy.strategy == RetryStrategy.LINEAR: delay = policy.base_delay * attempt elif policy.strategy == RetryStrategy.EXPONENTIAL: delay = policy.base_delay * (2 ** attempt) delay = min(delay, policy.max_delay) if policy.jitter: delay += random.uniform(0, delay * 0.1) return delay async def with_retry( func: callable, error_handler: callable = None, policies: SmartRetryPolicy = None ): """带智能重试的函数调用""" policies = policies or SmartRetryPolicy() last_error = None for attempt in range(10): # 安全上限 try: return await func() except Exception as e: last_error = e policy = policies.get_policy(e) if attempt >= policy.max_attempts: raise # 重试前操作(如修复 Prompt) if policy.pre_action == "repair_prompt": if error_handler: await error_handler(e, attempt) delay = policies.compute_delay(attempt, policy) logger.warning( f"Attempt {attempt+1} failed: {e}. " f"Retrying in {delay:.1f}s..." ) await asyncio.sleep(delay) raise last_error 2.2 格式错误自修复 class FormatRepairAgent: """格式错误自动修复""" REPAIR_STRATEGIES = [ "json_extract", # 从文本中提取 JSON "json_repair", # 修复常见 JSON 错误 "re_prompt", # 重新请求 LLM "structured_output", # 切换到结构化输出 ] async def repair(self, raw_output: str, expected_schema: dict) -> dict: """尝试修复格式错误""" for strategy in self.REPAIR_STRATEGIES: try: result = await getattr(self, f"_strategy_{strategy}")( raw_output, expected_schema ) if result is not None: logger.info(f"Format repaired using: {strategy}") return result except Exception: continue raise FormatRepairError( f"Could not repair output after all strategies" ) async def _strategy_json_extract(self, raw: str, schema: dict) -> dict | None: """从文本中提取 JSON""" # 尝试找到 JSON 块 import re patterns = [ r'```json\s*(.*?)\s*```', # Markdown JSON 块 r'```\s*(.*?)\s*```', # 通用代码块 r'\{[^{}]*\}', # 裸 JSON 对象 r'\[.*\]', # 裸 JSON 数组 ] for pattern in patterns: match = re.search(pattern, raw, re.DOTALL) if match: try: return json.loads(match.group(1)) except json.JSONDecodeError: continue return None async def _strategy_json_repair(self, raw: str, schema: dict) -> dict | None: """修复常见 JSON 错误""" repaired = raw # 修复尾随逗号 repaired = re.sub(r',\s*}', '}', repaired) repaired = re.sub(r',\s*]', ']', repaired) # 修复单引号 repaired = repaired.replace("'", '"') # 修复未引用的键 repaired = re.sub(r'(\w+):', r'"\1":', repaired) # 修复省略号 repaired = repaired.replace('...', 'null') try: return json.loads(repaired) except json.JSONDecodeError: return None async def _strategy_re_prompt(self, raw: str, schema: dict) -> dict | None: """使用 LLM 修复格式""" repair_prompt = f"""The following text was supposed to be valid JSON matching this schema: {json.dumps(schema, indent=2)} But it had format errors. Fix it and return ONLY valid JSON: Original text: {raw[:2000]} Return the corrected JSON:""" response = await llm.invoke( repair_prompt, response_format={"type": "json_object"}, temperature=0.0 ) try: return json.loads(response.content) except: return None async def _strategy_structured_output(self, raw: str, schema: dict) -> dict | None: """使用结构化输出 API 重新请求""" response = await llm.invoke( f"Convert this to structured data:\n{raw[:2000]}", response_format={"type": "json_schema", "json_schema": schema}, temperature=0.0 ) return json.loads(response.content) if response.content else None 三、中间层:工具错误恢复 3.1 工具执行框架 class ResilientToolExecutor: """带错误恢复的工具执行器""" def __init__(self): self.fallback_chains = {} # {tool_name: [backup_tool1, backup_tool2]} self.error_handlers = {} # {tool_name: handler} async def execute( self, tool_call: ToolCall, context: ExecutionContext ) -> ToolResult: tool_name = tool_call.tool primary = self._get_tool(tool_name) try: result = await with_retry( lambda: primary.execute(tool_call.args, context), policies=self._get_retry_policy(tool_name) ) return ToolResult(success=True, data=result) except Exception as e: # 尝试错误处理器 handler = self.error_handlers.get(tool_name) if handler: repaired_call = await handler(e, tool_call, context) if repaired_call: return await self.execute(repaired_call, context) # 尝试备用工具链 for backup_name in self.fallback_chains.get(tool_name, []): backup = self._get_tool(backup_name) try: adapted_args = await self._adapt_args( tool_call.args, primary, backup ) result = await backup.execute(adapted_args, context) return ToolResult( success=True, data=result, degraded=True, used_fallback=backup_name ) except Exception: continue return ToolResult( success=False, error=str(e), tool=tool_name ) # 注册备用工具链 executor = ResilientToolExecutor() executor.fallback_chains = { "web_search": ["bing_search", "duckduckgo_search"], "database_query": ["cache_lookup", "default_response"], "send_email": ["queue_email", "save_to_retry"], } # 注册错误处理器 async def search_error_handler( error: Exception, tool_call: ToolCall, context: ExecutionContext ) -> ToolCall | None: """搜索工具错误处理器""" if isinstance(error, TimeoutError): # 减少结果数量重试 modified_args = tool_call.args.copy() modified_args["max_results"] = 3 # 减少结果数量 return ToolCall(tool=tool_call.tool, args=modified_args) elif isinstance(error, RateLimitError): return None # 不重试,直接走备用 return None executor.error_handlers["web_search"] = search_error_handler 3.2 LLM 输出验证与自修复 class OutputValidator: """LLM 输出验证与自修复""" async def validate_and_repair( self, output: str, task: str, constraints: list[str] ) -> ValidatedOutput: # 1. 结构验证 structural = self._check_structure(output, task) if not structural.valid: repaired = await self._repair_structure(output, structural.errors) if repaired: output = repaired # 2. 约束验证 constraint_results = await self._check_constraints(output, constraints) violated = [r for r in constraint_results if not r.passed] if not violated: return ValidatedOutput(valid=True, output=output) # 3. 自修复:让 LLM 自己修正 if self._is_repairable(violated): repaired_output = await self._self_repair( output, task, violated ) if repaired_output: # 重新验证 recheck = await self._check_constraints(repaired_output, constraints) if all(r.passed for r in recheck): return ValidatedOutput( valid=True, output=repaired_output, repaired=True ) return ValidatedOutput( valid=False, output=output, violations=[r.constraint for r in violated] ) async def _self_repair( self, original_output: str, task: str, violations: list[ConstraintViolation] ) -> str | None: """让 LLM 自我修复输出""" violation_desc = "\n".join( f"- {v.constraint}: {v.detail}" for v in violations ) repair_prompt = f"""Your previous response has issues that need to be fixed. Original task: {task} Your response: {original_output[:3000]} Issues found: {violation_desc} Please fix these issues and provide the corrected response. Only output the corrected response, no explanations.""" try: response = await llm.invoke(repair_prompt, temperature=0.0) return response.content except: return None 四、高级层:自修复 Agent class SelfHealingAgent: """自修复 Agent:能识别错误并自主修复""" async def run_with_healing( self, task: str, max_heal_attempts: int = 3 ) -> str: for attempt in range(max_heal_attempts): try: result = await self._execute(task) # 自我验证 validation = await self._self_validate(task, result) if validation.confidence > 0.8: return result # 识别问题 diagnosis = await self._diagnose( task, result, validation ) # 生成修复方案 fix_plan = await self._generate_fix_plan(diagnosis) # 执行修复 task = await self._apply_fix(task, fix_plan) logger.info( f"Self-healing attempt {attempt+1}: " f"diagnosis={diagnosis.issue}, " f"fix={fix_plan.action}" ) except CircularError as e: # 检测到循环,换策略 task = await self._reframe_task(task, e) except MaxIterationsError: # 达到迭代上限,简化任务 task = await self._simplify_task(task) # 自修复失败,返回最佳结果 return result or "I was unable to complete this task." async def _diagnose( self, task: str, result: str, validation: Validation ) -> Diagnosis: """诊断输出问题""" diagnosis_prompt = f"""Analyze why the following result may be incorrect: Task: {task} Result: {result[:2000]} Validation concerns: {validation.concerns} Identify the most likely issue: - wrong_approach: Used wrong method to solve the problem - incomplete: Missing important information or steps - hallucination: Contains fabricated information - format_error: Output format doesn't match requirements - logical_error: Reasoning contains logical flaws - other: Specify Respond in JSON: {{"issue": "...", "detail": "...", "confidence": 0.0-1.0}}""" response = await llm.invoke(diagnosis_prompt, temperature=0.0) return Diagnosis(**json.loads(response.content)) async def _generate_fix_plan(self, diagnosis: Diagnosis) -> FixPlan: """生成修复计划""" FIX_STRATEGIES = { "wrong_approach": "Rethink the approach and try a different method", "incomplete": "Identify missing information and supplement it", "hallucination": "Verify all claims using tools and remove unverified ones", "format_error": "Reformat the output to match requirements", "logical_error": "Review the reasoning chain step by step", } action = FIX_STRATEGIES.get(diagnosis.issue, "Retry with more careful approach") return FixPlan( action=action, issue=diagnosis.issue, detail=diagnosis.detail, strategy="adjust_and_retry" ) 五、错误恢复编排 class ErrorRecoveryOrchestrator: """错误恢复编排器""" async def execute_with_recovery( self, task: str, agent: Agent ) -> RecoveryResult: recovery_layers = [ ("retry", RetryLayer()), ("format_repair", FormatRepairLayer()), ("tool_fallback", ToolFallbackLayer()), ("output_validation", OutputValidationLayer()), ("self_healing", SelfHealingLayer()), ("human_escalation", HumanEscalationLayer()), ] current_task = task current_output = None for layer_name, layer in recovery_layers: try: result = await layer.process( current_task, current_output, agent ) if result.success: return RecoveryResult( output=result.output, recovery_layers_used=[layer_name], success=True ) # 更新任务或输出,传给下一层 if result.modified_task: current_task = result.modified_task if result.partial_output: current_output = result.partial_output except Exception as e: logger.error(f"Recovery layer {layer_name} failed: {e}") continue return RecoveryResult( output=current_output or "Unable to complete task", success=False, recovery_layers_used=["all_failed"] ) 六、错误恢复 Checklist □ 错误分类体系明确(可重试/不可重试/需判断) □ 重试策略匹配错误类型(指数退避/线性/固定) □ 格式错误自动修复(JSON提取/修复/重新请求) □ 工具备用链配置(主工具失败→备用工具) □ LLM 输出验证+自我修复 □ 自修复 Agent 能诊断并修复自身错误 □ 循环检测器防止自修复死循环 □ 人工升级机制(自修复失败时触发) □ 错误恢复全链路追踪 □ 定期演练错误恢复流程 结语 错误恢复是 Agent 可靠性的最后一道防线。从简单的重试到复杂的自修复,每一层恢复策略都在增加系统的韧性。但要注意:自修复不是万能的——它消耗额外 Token、增加延迟、可能引入新错误。最佳策略是分层防御:基础错误用重试解决,格式错误用修复解决,逻辑错误用自修复解决,复杂错误交给人工。让 Agent 像人一样——犯错后能自我纠正,但也知道什么时候该求助。 加入讨论 这篇文章有姊妹讨论帖在硅基AGI论坛 — 全球首个碳基硅基认知交流平台。 ...

2026-06-28 · 7 min · 1399 words · 硅基 AGI 探索者
llm hallucination 2026 analysis

大模型幻觉问题 2026:根因分析与缓解策略

幻觉:LLM 最棘手的问题 大模型的幻觉(Hallucination)——自信地输出错误或虚构的信息——是 2026 年 LLM 生产应用中最大的痛点。Stanford 2026 年 AI Index 报告显示,即使是最先进的模型,在事实性问答中的幻觉率仍有 3-8%,在专业领域更是高达 15-20%。幻觉不是 bug,而是 LLM 生成式架构的特性——但我们可以系统性地缓解它。 一、幻觉的分类 1.1 幻觉类型体系 类型 描述 示例 严重性 事实性幻觉 输出与事实不符 “爱因斯坦出生于1890年”(实际1879) 高 来源幻觉 虚构引用/出处 编造不存在的论文 高 推理幻觉 推理链条中出错 正确前提但错误结论 中 时间幻觉 时间线混乱 “2024年奥运会在北京举办” 中 实体幻觉 虚构人/组织/产品 “微软CEO John Smith” 高 数值幻觉 数字/计算错误 “12×13=146”(实际156) 高 代码幻觉 API/函数不存在 调用不存在的库函数 中 自我幻觉 虚构自身能力 “我可以访问实时互联网” 低 1.2 幻觉产生的根因 ┌──────────────────────────────────────────┐ │ 幻觉根因分析 │ ├──────────────────────────────────────────┤ │ │ │ 1. 训练数据问题 │ │ ├── 训练数据中的错误信息 │ │ ├── 知识截止日期后的信息缺失 │ │ └── 长尾知识表示不足 │ │ │ │ 2. 模型架构问题 │ │ ├── 自回归生成无法回溯修正 │ │ ├── 注意力机制对事实关注不足 │ │ └── 知识存储与检索机制不完善 │ │ │ │ 3. 推理机制问题 │ │ ├── 缺乏事实核查机制 │ │ ├── 过度依赖模式匹配而非知识检索 │ │ └── 采样温度增加随机性 │ │ │ │ 4. 交互问题 │ │ ├── 模型倾向"回答"而非"承认不知道" │ │ ├── 用户问题中的错误前提引导 │ │ └── 上下文中的错误信息被采纳 │ │ │ └──────────────────────────────────────────┘ 二、幻觉检测方法 2.1 基于一致性的检测 class ConsistencyBasedDetector: """基于一致性的幻觉检测""" def __init__(self, llm_client, n_samples: int = 5): self.llm = llm_client self.n_samples = n_samples def detect(self, query: str, response: str) -> dict: """通过多次采样检测一致性""" # 1. 对同一问题生成多个回答 responses = [] for _ in range(self.n_samples): r = self.llm.generate(query, temperature=0.7) responses.append(r) # 2. 提取每个回答中的事实性陈述 all_claims = [] for r in responses: claims = self._extract_claims(r) all_claims.append(claims) # 3. 检查事实一致性 consistency_scores = self._check_consistency(all_claims) # 4. 原始回答的事实性评估 original_claims = self._extract_claims(response) original_scores = [] for claim in original_claims: # 该claim在其他回答中出现的比例 support_count = sum( 1 for claims in all_claims if self._claim_match(claim, claims) ) original_scores.append({ 'claim': claim, 'support_rate': support_count / len(all_claims), 'likely_hallucination': support_count / len(all_claims) < 0.4 }) hallucination_rate = sum( 1 for s in original_scores if s['likely_hallucination'] ) / len(original_scores) if original_scores else 0 return { 'hallucination_rate': hallucination_rate, 'claim_scores': original_scores, 'overall_consistency': np.mean(consistency_scores) } def _extract_claims(self, text: str) -> list: """提取文本中的事实性陈述""" prompt = f"""提取以下文本中的所有事实性陈述,每条一行: {text}""" result = self.llm.generate(prompt) return [line.strip() for line in result.split('\n') if line.strip()] def _claim_match(self, claim: str, other_claims: list) -> bool: """判断claim是否在other_claims中存在""" prompt = f"""判断以下陈述是否语义等价: 陈述A:{claim} 陈述B列表:{other_claims} 如果A与B中任一陈述语义等价,返回"是",否则"否"。""" return '是' in self.llm.generate(prompt) 2.2 基于来源验证的检测 class SourceVerificationDetector: """基于来源验证的幻觉检测""" def __init__(self, llm_client, search_client): self.llm = llm_client self.search = search_client def detect(self, response: str, sources: list = None) -> dict: """验证回答中的事实是否有来源支持""" # 1. 提取事实性陈述 claims = self._extract_claims(response) # 2. 对每个claim进行验证 results = [] for claim in claims: # 如果有现成来源,在来源中验证 if sources: verification = self._verify_in_sources(claim, sources) else: # 搜索验证 verification = self._verify_by_search(claim) results.append({ 'claim': claim, 'verdict': verification['verdict'], # supported | refuted | unverifiable 'evidence': verification['evidence'], 'confidence': verification['confidence'] }) hallucination_claims = [ r for r in results if r['verdict'] == 'refuted' ] unsupported_claims = [ r for r in results if r['verdict'] == 'unverifiable' ] return { 'total_claims': len(results), 'supported': len([r for r in results if r['verdict'] == 'supported']), 'refuted': len(hallucination_claims), 'unverifiable': len(unsupported_claims), 'hallucination_rate': len(hallucination_claims) / len(results), 'details': results } def _verify_in_sources(self, claim: str, sources: list) -> dict: """在给定来源中验证claim""" prompt = f"""判断以下陈述是否被来源文本支持。 陈述:{claim} 来源文本: {chr(10).join(f'[{i+1}] {s[:500]}' for i, s in enumerate(sources))} 判断: - "supported":来源支持该陈述 - "refuted":来源与该陈述矛盾 - "unverifiable":来源无法验证该陈述 返回判断和依据。""" result = self.llm.generate(prompt) if 'supported' in result.lower(): verdict = 'supported' elif 'refuted' in result.lower(): verdict = 'refuted' else: verdict = 'unverifiable' return { 'verdict': verdict, 'evidence': result, 'confidence': 0.8 } def _verify_by_search(self, claim: str) -> dict: """通过搜索验证claim""" search_results = self.search.search(claim, top_k=3) return self._verify_in_sources(claim, search_results) 2.3 基于模型置信度的检测 class ConfidenceBasedDetector: """基于模型内部置信度的幻觉检测""" def __init__(self, model): self.model = model def detect(self, prompt: str) -> dict: """通过logits分析检测幻觉""" # 1. 获取模型输出和logits with torch.no_grad(): outputs = self.model.generate( prompt, return_logits=True, output_scores=True ) # 2. 计算每个token的置信度 token_confidences = [] for i, (token, score) in enumerate( zip(outputs.tokens, outputs.scores) ): # softmax得到概率 probs = torch.softmax(score, dim=-1) token_prob = probs[token].item() token_confidences.append(token_prob) # 3. 分析低置信度区域 low_confidence_threshold = 0.5 low_conf_regions = [ {'position': i, 'token': outputs.tokens[i], 'confidence': conf} for i, conf in enumerate(token_confidences) if conf < low_confidence_threshold ] # 4. 计算整体置信度指标 avg_confidence = np.mean(token_confidences) min_confidence = np.min(token_confidences) return { 'avg_confidence': avg_confidence, 'min_confidence': min_confidence, 'low_confidence_tokens': len(low_conf_regions), 'hallucination_risk': 'high' if avg_confidence < 0.6 else 'medium' if avg_confidence < 0.8 else 'low', 'low_conf_regions': low_conf_regions[:10] # 前10个低置信度token } 三、幻觉缓解策略 3.1 RAG 增强 class RAGHallucinationMitigator: """RAG 增强缓解幻觉""" def __init__(self, llm_client, retrieval_client): self.llm = llm_client self.retrieval = retrieval_client def generate(self, query: str) -> dict: """RAG增强生成""" # 1. 检索相关文档 docs = self.retrieval.search(query, top_k=5) # 2. 构建RAG Prompt prompt = self._build_rag_prompt(query, docs) # 3. 生成回答 response = self.llm.generate(prompt) # 4. 后验证 verification = self._verify_response(response, docs) return { 'response': response, 'sources': docs, 'verification': verification, 'hallucination_risk': verification['hallucination_rate'] } def _build_rag_prompt(self, query: str, docs: list) -> str: return f"""请基于以下参考资料回答问题。如果资料中没有答案,请明确说明"根据现有资料无法回答"。 ## 参考资料 {self._format_docs(docs)} ## 回答规则 1. 只使用参考资料中的信息 2. 每个事实性陈述必须标注来源 [1], [2] 等 3. 如果资料中有矛盾信息,指出矛盾 4. 不确定的信息要标注"可能" 5. 资料中未涉及的信息不要编造 ## 问题 {query} ## 回答""" 3.2 自我验证机制 class SelfVerificationGenerator: """自我验证生成机制""" def __init__(self, llm_client): self.llm = llm_client def generate_with_verification(self, query: str) -> dict: """带自我验证的生成""" # Step 1: 初步生成 initial_response = self.llm.generate(query) # Step 2: 自我事实核查 fact_check = self._self_fact_check(query, initial_response) # Step 3: 如果发现幻觉,修正 if fact_check['has_issues']: corrected = self._correct_hallucinations( query, initial_response, fact_check ) else: corrected = initial_response # Step 4: 最终验证 final_check = self._final_verification(query, corrected) return { 'response': corrected, 'initial_response': initial_response, 'fact_check': fact_check, 'final_verification': final_check, 'corrections_made': fact_check['issues_found'] } def _self_fact_check(self, query: str, response: str) -> dict: """自我事实核查""" prompt = f"""请对你自己的回答进行事实核查。 问题:{query} 你的回答:{response} 请逐句检查: 1. 每个事实性陈述是否准确? 2. 是否有编造的信息? 3. 是否有不确定但表述为事实的内容? 4. 数字和日期是否正确? 对每个可能的问题,标注: - 陈述内容 - 问题类型(错误/编造/不确定) - 修正建议 如果一切正确,返回"无问题"。""" result = self.llm.generate(prompt) has_issues = '无问题' not in result issues = [] if has_issues: # 解析问题列表 issues = self._parse_issues(result) return { 'has_issues': has_issues, 'issues_found': issues, 'raw_check': result } def _correct_hallucinations(self, query, response, fact_check): """修正检测到的幻觉""" prompt = f"""请修正以下回答中的事实性错误。 原始问题:{query} 原始回答:{response} 检测到的问题:{fact_check['raw_check']} 修正规则: 1. 修正所有检测到的事实性错误 2. 不确定的信息添加"据我所知"等限定词 3. 无法确认的信息直接删除或标注"需查证" 4. 保持回答的连贯性和有用性 修正后的回答:""" return self.llm.generate(prompt) 3.3 多模型交叉验证 class CrossModelVerification: """多模型交叉验证""" def __init__(self, models: list): self.models = models # 多个不同厂商的模型 def generate_verified(self, query: str) -> dict: """多模型交叉验证生成""" # 1. 每个模型独立生成回答 responses = [] for model in self.models: r = model.generate(query, temperature=0.0) responses.append(r) # 2. 提取各模型的事实性陈述 all_claims = [] for r in responses: claims = self._extract_claims(r) all_claims.append(set(claims)) # 3. 找出共识陈述和分歧陈述 consensus = set.intersection(*all_claims) if all_claims else set() all_claims_set = set.union(*all_claims) if all_claims else set() disputed = all_claims_set - consensus # 4. 只保留共识陈述 verified_response = self._build_response_from_consensus( query, consensus, responses[0] # 以第一个模型的回答为模板 ) return { 'response': verified_response, 'model_count': len(self.models), 'consensus_claims': len(consensus), 'disputed_claims': len(disputed), 'agreement_rate': len(consensus) / max(len(all_claims_set), 1) } 四、幻觉评估基准 4.1 评估指标 指标 计算方式 目标 事实准确率 事实正确陈述/总事实陈述 > 95% 幻觉率 虚构陈述/总事实陈述 < 5% 来源准确率 正确引用/总引用 > 90% 拒绝准确率 正确拒绝不该回答的问题比例 > 95% 过度拒绝率 错误拒绝合理问题的比例 < 5% 自校准ECE 预期校准误差 < 0.1 4.2 评估数据集构建 class HallucinationEvalDataset: """幻觉评估数据集""" CATEGORIES = { 'factual_qa': { 'description': '事实性问答', 'examples': [ {'q': '中国最长的河流是?', 'a': '长江', 'type': 'fact'}, {'q': '光速是多少?', 'a': '约3×10^8 m/s', 'type': 'fact'}, ] }, 'unanswerable': { 'description': '无法回答的问题(测试是否承认不知道)', 'examples': [ {'q': '2028年的奥斯卡最佳影片是哪部?', 'a': '无法回答', 'type': 'refuse'}, {'q': '张三的手机号码是多少?', 'a': '无法回答', 'type': 'refuse'}, ] }, 'false_premise': { 'description': '错误前提问题', 'examples': [ {'q': '林黛玉倒拔垂杨柳的故事告诉我们什么?', 'a': '指出前提错误', 'type': 'correct'}, ] }, 'multi_hop': { 'description': '多跳推理(每跳都可能出错)', 'examples': [ {'q': '相对论提出者出生国家的首都人口是多少?', 'a': '需推理:爱因斯坦→德国→柏林→人口', 'type': 'reasoning'} ] } } 五、2026 前沿方向 检索增强自我修正:模型生成后自动检索验证并修正 事实性训练:用事实核查数据做偏好优化 知识编辑:直接在模型权重中修正错误知识 不确定性量化:模型输出附带校准的不确定性分数 因果推理增强:用因果推理替代模式匹配 工具增强事实性:自动调用搜索引擎、计算器等工具 结语 幻觉是 LLM 的"原罪"——源于其自回归生成架构和统计学习本质。完全消除幻觉在当前技术范式下不可能,但通过 RAG、自我验证、交叉验证等手段,可以将其控制在可接受范围内。 ...

2026-06-28 · 6 min · 1115 words · 硅基 AGI 探索者
鲁ICP备2026018361号