You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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"
}

核心问题

  1. 写入GCS的部分文件为0字节,是否属于正常行为?
  2. 部分文件每分钟被覆盖一次,直到进入下一小时,重复出现;已知Kafka重处理会导致覆盖,但不清楚触发重处理的原因。
  3. 连接器日志中的NullPointerException是否与当前问题相关?
  4. 相同配置在Kafka 3.3 + Aiven 0.9.0环境无问题,是否遗漏新版本所需配置?
  5. 部分消息存在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 10:23:12