题干与适用场景
商品数据库是搜索系统的唯一事实源。请设计一个近实时 CDC 管道,把新增、更新和删除可靠地同步到搜索索引;要求支持初始全量、增量追平、重复投递、消费中断、Schema 变化和索引重建,并说明如何证明没有漏数据或旧事件覆盖新事件。
这道题适合后端、数据基础设施、搜索平台和系统设计面试。题目假设数据库可以提供提交顺序或等价的日志位置,搜索索引是可重建的派生系统;不把 Kafka、Debezium 或 Elasticsearch 当成必选答案,候选人需要先说清楚语义与失败边界。
面试官考察点
面试官会看候选人能否把“数据库写成功”和“索引最终可见”拆成两个可监控的阶段。强回答会先定义每条事件的主键、操作类型、事务或日志位置和 Schema 版本,再选择日志型 CDC 而不是依赖轮询时间戳;随后处理快照与增量重叠、至少一次重放、按键顺序、删除墓碑和索引切换。只画数据库、消息队列和搜索框,却没有 checkpoint、回放和校验规则,无法证明可靠性。
回答前需要澄清的问题
- 新鲜度目标是什么? 是提交后 5 秒内可搜索,还是允许分钟级延迟?这决定缓冲、告警和降级预算。
- 必须保留什么顺序? 通常只要求同一商品按源端提交顺序应用,并不要求跨商品全局排序;若商品和库存跨表有一致性约束,需要额外设计聚合事件。
- 删除是硬删除还是软删除? 硬删除需要可靠墓碑或删除事件;软删除则要把可见性条件编码进索引文档。
- 全量期间允许写入吗? 若允许,就必须定义快照位置,并让该位置之后的增量补偿全量读取期间的变化。
- Schema 如何演进? 新字段能否被旧消费者忽略,字段删除是否需要双写、重建索引或版本化映射?
- 重建是否要求零停机? 若要求,使用新索引别名切换,并为旧消费者保留回放起点。
30 秒回答框架
“我会把主数据库当唯一事实源,优先从提交日志捕获已提交的 insert、update、delete,并给事件带主键、操作、源端 LSN、事务 ID、Schema 版本和前后值。首次同步从一个一致性快照开始,同时记录快照对应的日志位置;快照之后的事件进入同一条可重放流,消费者按商品键幂等写入索引。
事件处理采用至少一次交付,checkpoint 只在索引写成功后推进,重复事件用主键加版本或 LSN 条件写挡住。消费者中断后从 checkpoint 重放;积压、最老事件年龄、源端 slot 保留量和索引与数据库抽样校验都做成指标。重建时把同一事件流写入新索引,追平后再原子切 alias,并保留回放窗口。”
分步骤深入解答
第一步:定义事件契约和捕获边界
日志型 CDC 读取数据库提交日志,能看到事务提交后的变化,并保留源端顺序或日志位置。以 PostgreSQL 为例,logical decoding 从 WAL 提取变更,replication slot 表示可以按源端顺序重放的变更流;slot 会持有所需 WAL,因此必须监控保留量,避免连接器停滞拖满主库磁盘。
事件至少包含 entityid、operation、sourceposition、transactionid、schemaversion、before 和 after。把 source_position 当作审计和去重依据,而不是把消费者收到消息的时间当作业务顺序。一个事务影响多个商品时,要决定索引是否允许逐条可见,还是先在流中按事务边界聚合。
第二步:用快照位置连接全量和增量
全量与增量最危险的窗口是:快照读到旧值之后,流又收到同一主键的新值。启动快照时记录日志位置 P0;快照产生的文档只代表开始时的状态,P0 之后的事件必须继续保留并在快照结果后应用。
P0 = captureSourcePosition()
startStreaming(after=P0)
for row in consistentSnapshot():
indexUpsert(row, version=P0)
for event in stream:
if event.position > indexedVersion[event.key]:
applyIdempotently(event)
checkpoint(event.position) # only after index write succeeds真实连接器可能用快照窗口、主键分块和缓冲区处理 READ 与 UPDATE 的碰撞;面试中要说明这是为了避免旧的快照行覆盖已经提交的新事件,而不是简单地“快照完成后再开流”。
第三步:设计幂等、顺序和失败恢复
消费者按 entity_id 分区,让同一商品的事件保持源端顺序;不同商品可以并行。索引写入使用条件版本、外部版本号或带版本的文档,只有事件位置更新时才覆盖。DELETE 写入墓碑或删除操作,并保留足够的版本信息,防止迟到 UPDATE 把已删除文档复活。
checkpoint 表示“该事件的副作用已成功持久化”,不能在拉到消息或发送 HTTP 请求后提前提交。进程在索引写成功、checkpoint 写入前崩溃时会重复投递,所以目标写入必须幂等;如果 checkpoint 领先于索引,则会产生漏数据,必须把两者放在同一可验证的提交边界,或使用可重放的索引任务和对账修复。
第四步:处理重放、Schema 和重建
每个消费者使用独立的 slot 或等价进度,避免多个消费者争抢同一单次消费游标。重放前冻结或标记目标索引的版本策略,限制回放范围,并让旧事件只能写入更旧版本。Schema 事件要有兼容规则:新增可选字段通常允许旧消费者忽略;删除或改类型可能需要新事件版本、双读双写或重新索引。
重建不应直接清空线上索引。创建新索引并从快照位置开始回放,直到新索引的已应用位置达到切换门槛;随后原子切换 alias,并继续消费同一流。若切换失败,保留旧 alias 和新索引的进度,修复后继续追平,不重新猜测起点。
第五步:用指标和对账证明可靠性
核心指标包括 CDC 读取延迟、分区积压、最老事件年龄、slot WAL 保留量、每个消费者 checkpoint、索引写入失败率、重试和死信数量,以及数据库与索引按主键抽样比较的版本差。对删除尤其要统计墓碑处理与残留文档。
验证应故意停止消费者、重复投递、打乱不同分区的消息、在快照期间更新同一商品、发送迟到删除、改变 Schema 并执行主库故障切换。对账工具要能从数据库重查当前版本、从事件流重放到某个位置,并输出最小不一致样本。没有可重放位置和样本,单看“队列为空”不能证明没有漏数据。
高质量示范回答
“我先定义契约:数据库是事实源,索引是可重建副本;目标是提交后 5 秒内可见,同一商品按源端提交顺序处理,跨商品不承诺全局排序。事件携带主键、insert/update/delete、LSN、事务 ID、Schema 版本和前后值。
我会用日志型 CDC。启动一致性快照时记下位置 P0,然后从 P0 之后的流继续消费。快照 READ 与流中的 UPDATE 可能碰撞,所以要用快照窗口或等价的主键版本去重,不能让旧 READ 覆盖新 UPDATE。消费者按商品键分区,索引使用外部版本或条件写;DELETE 要保留墓碑版本,避免迟到 UPDATE 复活。
checkpoint 只有在索引副作用成功后推进,崩溃会导致至少一次重放,因此写入必须幂等。连接器停机时我监控 slot 保留的 WAL,恢复后从最后安全位置继续。重建时写新索引并从同一位置回放,追平后原子切 alias。
验收不只看队列长度。我会注入快照期间更新、重复和乱序、迟到删除、消费者崩溃、Schema 变更与主库切换,再按主键比较数据库版本、索引版本和 checkpoint。重点指标是最老事件年龄、WAL 保留、索引延迟、死信和不一致样本;出现缺口时能从保存的日志位置重放和修复。”
常见错误
- 用更新时间轮询代替可靠 CDC → 时钟精度、回拨和长事务会漏掉变化 → 读取提交日志或明确可证明的游标。
- 快照和增量各自启动却没有共同位置 → 快照旧值可能覆盖新事件 → 记录 P0,并处理 READ/UPDATE 碰撞。
- checkpoint 在发送索引请求后立即推进 → 崩溃窗口会造成漏数据 → 只在副作用可验证成功后推进。
- 假设至少一次等于 exactly-once → 重复投递仍会发生 → 用版本条件写和幂等删除。
- 只按消息到达时间排序 → 网络重试会打乱源端顺序 → 按主键分区并使用源端 LSN 或版本。
- 删除直接从索引移除且不留版本 → 迟到更新会让文档复活 → 保存墓碑或删除版本。
- 多个消费者共享一个 replication slot → 变化可能被其中一个独占,其他消费者静默缺数据 → 每个独立消费者使用独立 slot 或明确的广播层。
- 重建时清空线上索引 → 回放失败会造成大面积不可搜索 → 新索引追平后原子切换 alias。
- 只看队列为空 → 可能已跳过事件或索引写入失败 → 做位置、版本和主库抽样对账。
追问及应对
追问一:如果一个事务更新了商品和库存,索引能逐条可见吗?
可以,但要明确业务允许短暂中间态;如果搜索结果必须同时反映两者,就要让事件携带事务边界并在消费者聚合后一次更新文档,或把可搜索视图建成一个已提交投影。不要用跨索引写入的“几乎同时”冒充原子性。
追问二:连接器停机很久导致 slot 占满 WAL,怎么止损?
先保护主库:告警并限制继续写入或临时降级消费者,确认 slot 是否仍有可用起点。若 slot 已失效,不能随意创建新 slot 继续跑,因为中间 LSN 可能丢失;应从备份或全量快照重建,并用对账证明缺口已补齐。
追问三:事件流只有分区级顺序,如何处理跨商品的搜索排序?
搜索排序读取的是索引当前版本,不应假设全局事件顺序。若排序字段需要全局一致时间,使用源端提交时间加版本规则,或由独立聚合器生成稳定排序键,并明确允许的暂时偏差。跨分区强制总序会牺牲吞吐,只有业务确实需要才采用。
追问四:Schema 删除字段时,旧索引怎么办?
先发布能同时读新旧事件的消费者,再停止产生旧字段,确认积压和重放窗口已清空,最后重建或迁移索引映射。若字段改变影响分析或权限语义,不能只删 JSON 字段;应保留版本化事件和回滚路径。