Kafka Aiven GCS Sink Connector空文件及重复覆盖问题排查
环境信息
- Aiven版本:
0.13.0 - Kafka版本:
3.5.1
连接器配置
{ "connector.class": "io.aiven.kafka.connect.gcs.GcsSinkConnector", "tasks.max": "4", "topics": "loan_logs", "gcs.credentials.path": "<credspath>", "gcs.bucket.name": "datalake-raw", "file.name.prefix": "lsm/", "file.name.timestamp.timezone": "Asia/Jakarta", "file.name.template": "{{topic}}/sink_date={{timestamp:unit=yyyy}}{{timestamp:unit=MM}}{{timestamp:unit=dd}}/sink_hour={{timestamp:unit=HH}}/p{{partition:padding=false}}-{{start_offset:padding=true}}.snappy.parquet", "file.compression.type": "snappy", "format.output.fields": "key,value,offset,timestamp,headers", "format.output.type": "parquet", "format.output.envelope": "true", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter.schemas.enable": "false", "behavior.on.null.values": "ignore" }
核心问题
- 写入GCS的部分文件为0字节,是否属于正常行为?
- 部分文件每分钟被覆盖一次,直到进入下一小时,重复出现;已知Kafka重处理会导致覆盖,但不清楚触发重处理的原因。
- 连接器日志中的
NullPointerException是否与当前问题相关? - 相同配置在Kafka 3.3 + Aiven 0.9.0环境无问题,是否遗漏新版本所需配置?
- 部分消息存在NULL key,是否与此问题相关?
日志报错
2024-02-09 04:06:01,665 ERROR [de-lsm-gcs-connector|task-1] WorkerSinkTask{id=de-lsm-gcs-connector-1} Offset commit failed, rewinding to last committed offsets (org.apache.kafka.connect.runtime.WorkerSinkTask) [task-thread-de-lsm-gcs-connector-1] org.apache.kafka.connect.errors.ConnectException: java.lang.NullPointerException at io.aiven.kafka.connect.gcs.GcsSinkTask.flushFile(GcsSinkTask.java:131) at java.base/java.util.HashMap.forEach(HashMap.java:1421) at java.base/java.util.Collections$UnmodifiableMap.forEach(Collections.java:1553) at io.aiven.kafka.connect.gcs.GcsSinkTask.flush(GcsSinkTask.java:114) at org.apache.kafka.connect.sink.SinkTask.preCommit(SinkTask.java:125) at org.apache.kafka.connect.runtime.WorkerSinkTask.commitOffsets(WorkerSinkTask.java:407) at org.apache.kafka.connect.runtime.WorkerSinkTask.commitOffsets(WorkerSinkTask.java:377) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:221) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:206) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:202) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:257) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:177) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833) Caused by: java.lang.NullPointerException ... 2024-02-09 04:06:01,665 ERROR [de-lsm-gcs-connector|task-1] WorkerSinkTask{id=de-lsm-gcs-connector-1} Commit of offsets threw an unexpected exception for sequence number 124: null (org.apache.kafka.connect.runtime.WorkerSinkTask) [task-thread-de-lsm-gcs-connector-1] org.apache.kafka.connect.errors.ConnectException: java.lang.NullPointerException at io.aiven.kafka.connect.gcs.GcsSinkTask.flushFile(GcsSinkTask.java:131) at java.base/java.util.HashMap.forEach(HashMap.java:1421) at java.base/java.util.Collections$UnmodifiableMap.forEach(Collections.java:1553) at io.aiven.kafka.connect.gcs.GcsSinkTask.flush(GcsSinkTask.java:114) at org.apache.kafka.connect.sink.SinkTask.preCommit(SinkTask.java:125) at org.apache.kafka.connect.runtime.WorkerSinkTask.commitOffsets(WorkerSinkTask.java:407) at org.apache.kafka.connect.runtime.WorkerSinkTask.commitOffsets(WorkerSinkTask.java:377) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:221) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:206) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:202) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:257) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:177) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833) Caused by: java.lang.NullPointerException ...
问题解答
1. 0字节文件是否正常?
不属于正常行为。空文件通常是任务在写入/刷新文件时发生异常(如日志中的NPE),导致文件被创建但未写入任何数据就中断;或是flush阶段失败,留下未完成的空文件。
2. 文件重复覆盖与消息重处理的原因
日志中的Offset commit failed, rewinding to last committed offsets是核心诱因:当偏移量提交失败时,Kafka Connect会回滚到上一次提交的偏移量,重新处理这段区间的消息,进而导致文件被重复覆盖。而偏移量提交失败的直接原因就是NullPointerException——它发生在GcsSinkTask.flushFile方法中,说明任务在将文件刷新到GCS时出现空指针错误,触发了重处理循环。
3. NPE错误与当前问题的关联
完全相关。NPE导致flush操作失败,进而引发偏移量提交失败,触发消息重处理,最终导致文件重复覆盖;同时,flush失败也会直接产生空文件。
4. 新版本是否遗漏配置?
是的。新版本的Aiven GCS连接器或Kafka Connect可能对NULL值的处理逻辑有变更,旧版本的兼容配置不再适用。例如,StringConverter在新版本中可能需要显式配置NULL key的处理策略,而旧版本(0.9.0)默认做了兼容处理。
5. NULL key的影响
NULL key是大概率的问题根源。当使用StringConverter且key.converter.schemas.enable=false时,NULL key会导致转换器返回null值,而新版本的GCS连接器未对该null值做防御性检查,触发NPE。旧版本的连接器可能有专门的兼容逻辑,因此未出现问题。
修复建议
- 处理NULL key:
- 添加
key.converter.null.handling.mode=ignore或key.converter.null.handling.mode=convertToDefault(具体参数参考StringConverter文档),避免NULL key导致转换错误。 - 改用
org.apache.kafka.connect.json.JsonConverter,它对NULL值的处理更友好,同时保持key.converter.schemas.enable=false。
- 添加
- 升级连接器版本:检查Aiven GCS连接器的最新版本,确认是否有修复NULL key导致NPE的bug,升级到最新稳定版。
- 调整文件滚动策略:添加
file.roll.interval.ms或file.roll.size.bytes配置,避免文件长时间处于未完成状态,减少重处理时的覆盖频率。 - 验证GCS权限:确保连接器拥有GCS的写入、更新权限,排除因权限问题导致flush失败的可能。
内容的提问来源于stack exchange,提问作者suisen

