Kafka-Snowflake连接器写入Snowflake记录数少于源端的排查方法咨询
前置:统一统计口径
两边计数前必须满足两个前提,否则所有对数结果都没有参考价值:
- 固定统计窗口:等源端停止向对应Topic生产数据,且Connector消费Lag清零、Snowflake管道无待加载文件后再开始统计,边写边对数必然存在偏差
- 统一统计维度:按相同的Topic、分区范围统计,排除Kafka侧已被保留策略清理的消息、Snowflake侧已被删除/归档的记录
Kafka侧核查机制与操作步骤
Kafka侧的可信记录数只认分区位点差计算的结果,分布式部署场景下消费位点统一存储在Kafka内置的__consumer_offsets主题中,单Connector Worker的日志计数、第三方监控面板的展示值大多存在抽样误差,不能作为对账依据。
- 确认消费停滞状态
执行以下命令查看Snowflake Connector对应消费组的状态,等所有分区的LAG值为0且5分钟内无波动,再进入后续计数步骤:
kafka-consumer-groups.sh --bootstrap-server <Kafka集群Broker地址> --describe --group <Snowflake Connector配置的消费组ID>
如果存在分区Lag长期不为0,先排查对应消费Worker的报错,解决消费阻塞问题再继续。
2. 计算Kafka侧留存的总有效记录数
分别查询每个分区的最早可用位点、最新生产位点:
# 查询各分区最早可用位点(小于这个位点的消息已被清理,无法被消费) kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list <Kafka集群Broker地址> --topic <对账Topic名> --time -2 # 查询各分区最新生产位点 kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list <Kafka集群Broker地址> --topic <对账Topic名> --time -1
单分区有效记录数 = 分区最新位点值 - 分区最早可用位点值,所有分区计算结果求和,就是Kafka侧当前可被Connector消费的总记录数。
如果某分区最早位点大于0,说明该分区已有历史消息因保留过期、被删除,这部分消息Connector永远无法读取,是数差的高频原因
3. 校验消费位点合理性
对比消费组的已提交位点和分区最新位点:
- 如果某分区已提交位点 > 分区最新生产位点,说明存在超前提交,对应区间的消息会被直接跳过,不会写入Snowflake
- 如果某分区已提交位点长期无推进,说明消费阻塞,优先排查Worker日志里的序列化错误、网络连通性问题
Snowflake侧核查机制与操作步骤
Snowflake侧的可信记录数必须排除重复写入、未加载、写入失败的记录,直接对目标表执行count(*)的结果没有参考价值——Snowflake Kafka Connector默认是至少一次投递语义,重试场景下必然存在重复数据。
- 确认管道加载完成状态
在Snowflake中执行以下命令,确认管道无待加载文件:
SELECT SYSTEM$PIPE_STATUS('<目标管道全路径,格式为 库名.模式名.管道名>');
返回结果中pendingFileCount为0、lastReceivedMessageTimestamp与当前时间差小于配置的刷盘间隔时,说明所有推送至Snowflake Stage层的文件都已完成加载。
2. 统计目标表有效记录数
用Kafka自带的分区+offset作为唯一键去重统计,排除重复写入的记录:
SELECT COUNT(DISTINCT RECORD_METADATA:partition::VARCHAR || '-' || RECORD_METADATA:offset::VARCHAR) AS valid_cnt FROM <目标写入表名>;
- 统计写入失败的丢弃记录数
Connector写入时解析失败、类型不匹配、超长的记录,会被写入默认后缀为_RECORD的错误表中,这部分记录不会进入正式表:
SELECT COUNT(*) AS failed_cnt FROM <错误表全路径名>;
- 校验表操作记录
查询Query History,确认统计时间窗口内目标表没有执行过DELETE、TRUNCATE、数据归档操作,避免人为操作导致的记录缺失。
数差根因定位逻辑
- 若 Kafka侧总记录数 = Snowflake有效记录数 + 错误表记录数:数差来自写入失败的脏数据,逐行处理错误表中的脏记录即可补全数据
- 若 Kafka侧总记录数 > Snowflake有效记录数 + 错误表记录数:优先核查两项配置:一是Connector是否配置了
behavior.on.null.values=drop,空值消息会被直接跳过不写入;二是Kafka消息保留时长是否短于Connector消费延迟,导致消息还没被消费就被清理 - 若排查上述配置后仍存在差值:检查分布式Connector集群的重平衡日志,确认是否存在Worker宕机、网络闪断导致重平衡时位点提交异常,跳过了部分消息。这种场景可以重置消费组位点到差值区间的起始位置,重新消费做一次全量校验。
内容的提问来源于stack exchange,提问作者Austin Jackson

