Kafka Connect S3 Sink按记录数或时间自动刷新配置问题排查
Kafka Connect S3 Sink 不按时刷新问题排查与修复
你的配置中遗漏了强制时间触发刷新的核心参数,同时部分时间分区参数的配合逻辑存在冲突,导致文件打开后无法按时刷写到S3,具体分析和修复如下:
1. 必须添加 flush.interval.ms 参数
你仅配置了 flush.size(按记录数触发刷新),但未设置 flush.interval.ms——这个参数的作用是当未达到flush.size的记录数时,每隔指定时间强制将缓冲区数据刷写到S3。默认情况下该参数未启用,因此只有当记录数凑够1000条时才会触发刷新,即使过了30秒也不会主动执行。
在连接器的config中添加:
"flush.interval.ms": "30000"
2. 调整 partition.duration.ms 与 path.format 的匹配
你当前设置 partition.duration.ms=30000(30秒),但 path.format 只定义到小时层级('hour'=HH),这会导致同一小时内生成的所有30秒分区都共享同一个目录,引发分区逻辑混乱,进而影响文件的滚动和刷新行为。
如果你的需求是按小时分区,建议将 partition.duration.ms 调整为与path.format匹配的1小时:
"partition.duration.ms": "3600000"
若确实需要按30秒分区,需修改path.format添加分钟/秒层级:
"path.format": "'year'=YYYY/'month'=MM/'day'=dd/'hour'=HH/'minute'=mm/'second'=ss"
3. 关于 rotate.interval.ms 的说明
rotate.interval.ms 控制的是同一分区内文件的滚动间隔,即到达该时间后,连接器会关闭当前文件并创建新文件。它需要与 flush.interval.ms 配合使用:flush.interval.ms 负责刷新数据到S3,rotate.interval.ms 负责切换文件。建议将两者设置为相同或成倍数的时间值,避免逻辑冲突。
修正后的完整配置示例
{ "name": "my-s3-sink", "config": { "connector.class": "io.confluent.connect.s3.S3SinkConnector", "tasks.max": "1", "topics": "my-topic", "topics.dir": "my-dir", "behavior.on.null.values": "ignore", "s3.region": "us-west-1", "s3.bucket.name": "by-bucket", "s3.part.size": "5242880", "aws.access.key.id": "${AWS_ACCESS_KEY}", "aws.secret.access.key": "${AWS_SECRET_KEY}", "storage.class": "io.confluent.connect.s3.storage.S3Storage", "flush.size": "1000", "flush.interval.ms": "30000", // 新增:每30秒强制刷新 "format.class": "io.confluent.connect.s3.format.avro.AvroFormat", "schema.compatibility": "NONE", "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner", "path.format": "'year'=YYYY/'month'=MM/'day'=dd/'hour'=HH", "partition.duration.ms": "3600000", // 调整为1小时分区 "rotate.interval.ms": "30000", "locale": "en-US", "timezone": "UTC" } }
内容的提问来源于stack exchange,提问作者Mr T.
相关产品推荐
相关产品推荐

