Confluent S3 Sink Connector设大参数仍每小时生成百万小文件求助
问题背景
尽管已设置较大的flush.size和rotate.interval.ms参数,我司的Confluent S3 Sink Connector每小时仍生成约100万个仅含少量记录(约3KB)的小文件,而非预期的大文件。
根因推测
两个归属生产团队的Debezium CDC连接器从分片数据库向同一主题写入数据,数据库字段完全相同但字段顺序不同:
Debezium Instance A: {"id": 1, "name": "test", "ts": 123456789} Debezium Instance B: {"ts": 123456789, "name": "test", "id": 1}
这导致Schema Registry中该主题存在多个schema版本,S3连接器会将每个schema版本视为独立写入单元,从而无视flush.size等阈值提前触发文件轮转。
当前生产环境配置
{ "connector.class": "io.confluent.connect.s3.S3SinkConnector", "storage.class": "io.confluent.connect.s3.storage.S3Storage", "s3.region": "us-west-2", "s3.bucket.name": "my-data-lake", "topics": "my-cdc-topic", "topics.dir": "warehouse", "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat", "parquet.codec": "gzip", "flush.size": "100000", "rotate.interval.ms": "3600000", "tasks.max": "20", "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner", "timestamp.extractor": "RecordField", "timestamp.field": "updated_at", "path.format": "YYYY/MM/dd/HH", "partition.duration.ms": "3600000", "schema.compatibility": "FULL", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://schema-registry:8081", "value.converter.schemas.enable": "false", "errors.tolerance": "all" }
已尝试的无效配置
除自定义SMT外,尝试过以下参数组合但问题未解决:
{ "value.converter.use.latest.version":"true", "value.converter.schemas.enable":"true", "enhanced.avro.schema.support":"true" }
解决方案建议
1. 统一Debezium CDC输出的字段顺序
修改两个Debezium实例的配置,强制输出相同的字段顺序:
- 若分片数据库表结构字段顺序不一致,调整表结构对齐字段顺序
- 配置Debezium的
column.whitelist参数,指定固定的字段顺序列表,确保两个实例生成的JSON结构字段顺序完全一致,避免Schema Registry生成多版本schema
2. 使用SchemaNormalizer SMT标准化记录结构
添加SMT到S3 Sink Connector配置中,强制将记录字段顺序统一为指定schema的顺序:
"transforms": "normalizeSchema", "transforms.normalizeSchema.type": "io.confluent.connect.transforms.SchemaNormalizer", "transforms.normalizeSchema.schema.name": "my-cdc-topic-value", "transforms.normalizeSchema.use.latest.version": "true"
该SMT会将输入记录的字段重新排序,匹配指定的schema结构,确保同一分区内的记录使用同一schema,避免S3连接器为不同schema版本创建单独写入器。
3. 调整Schema兼容性级别
将当前schema.compatibility: FULL调整为BACKWARD或FORWARD,允许字段顺序变化而不生成新的schema版本(需确保业务场景兼容字段顺序变化):
"schema.compatibility": "BACKWARD"
4. 启用Avro字段顺序忽略配置
为Avro Converter添加字段顺序忽略参数,让Avro忽略字段顺序差异,使用同一schema版本处理记录:
"value.converter.avro.field.order": "ignore"
此配置需配合schema.compatibility调整使用,确保Schema Registry接受字段顺序变化的记录而不生成新schema。
5. 注册统一标准Schema
若已存在多版本schema,可在Schema Registry中注册统一的标准schema,强制Debezium和S3 Connector使用该版本:
- 导出当前主题的最新schema,调整为固定字段顺序
- 在Schema Registry中注册该schema为新的版本(操作前需评估影响)
- 配置Debezium使用该schema序列化,S3 Connector使用该schema反序列化
内容的提问来源于stack exchange,提问作者Priyanshu sharma M

