Kafka Connect同步Protobuf数据到PostgreSQL字段丢失问题排查
Kafka Connect同步Protobuf数据到PostgreSQL时大量字段丢失
配置Kafka Connect将Kafka主题中的Protobuf格式数据同步至PostgreSQL数据库,同步可正常运行,但半数字段完全丢失,无任何报错。数据库中仅存在value、isSomething、physicalType和key字段,丢失字段涵盖string、int64类型及需转换为字符串的自定义类型。
连接器配置
{ "schema.registry.url": "http://schema-registry:8081", "name": "JdbcSinkConnectorConnector_0", "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schema.registry.url": "http://schema-registry:8081", "value.converter": "io.confluent.connect.protobuf.ProtobufConverter", "value.converter.schema.registry.url": "http://schema-registry:8081", "value.converter.schemas.enable": true, "topics": "schema.data.normalized", "confluent.controlcenter.schema.registry.url": "http://schema-registry:8081", "connection.url": "jdbc:postgresql://db:5432/db", "connection.user": "postgres", "connection.password": "a", "dialect.name": "PostgreSqlDatabaseDialect", "pk.mode": "record_key", "pk.fields": "key", "auto.create": true, "auto.evolve": true, "table.name.format": "app.TestTable", "fields.whitelist": "fieldA,created,unitSymbol,unitMultiplier,timeStamp,value,validForStart,validForEnd,isSomething,physicalType", "transforms": "RenameField", "transforms.RenameField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.RenameField.renames": "created:created_timestamp" }
Protobuf Schema
syntax = "proto3"; package xxx; import "identified_object.proto"; import "unit_symbol.proto"; import "unit_multiplier.proto"; option java_multiple_files = true; message XxxMessage { .IdentifiedObject id = 1; .IdentifiedObject environmental_data_provider = 2; string fieldA = 3; int64 created = 4; .UnitSymbol unit_symbol = 5; .UnitMultiplier unit_multiplier = 6; int64 time_stamp = 7; float value = 8; int64 valid_for_start = 9; int64 valid_for_end = 10; bool isSomething = 11; string physicalType = 12; }
示例消息
{ "id": { "aliasName": "", "description": "", "mRid": "8f90016b-7f2c-4172-a52a-bd6caae95020", "name": "" }, "environmentalDataProvider": { "aliasName": "", "description": "", "mRid": "408851cf-d0d5-43fe-a8cf-e234583aa7ae", "name": "...." }, "fieldA": "rb771347-3691-4920-88af-b2b8caffdea1", "created": "1664469000000", "unitSymbol": "UNIT_SYMBOL_METER_Per_SEC", "unitMultiplier": "UNIT_MULTIPLIER_UNSPECIFIED", "timeStamp": "1664857800000", "value": 4.2018013, "validForStart": "1664857800000", "validForEnd": "1664858700000", "isSomething": true, "physicalType": "typeA" }
问题排查与解决
修正字段命名匹配问题
Protobuf Schema使用下划线命名(如time_stamp、valid_for_start),但fields.whitelist中用了驼峰命名,导致转换器无法匹配到实际字段,最终被过滤。需将白名单字段改为Schema中的下划线格式:fields.whitelist: "fieldA,created,unit_symbol,unit_multiplier,time_stamp,value,valid_for_start,valid_for_end,isSomething,physicalType"保持
RenameField配置中created:created_timestamp的映射,因为原字段名是created,与Schema一致。处理自定义枚举类型
unit_symbol和unit_multiplier是自定义枚举,默认转换器不会自动转为字符串。添加以下配置让枚举直接转为字符串值:value.converter.enum.as.string: true验证Schema Registry中的Schema版本
确认Schema Registry中存储的XxxMessageSchema与本地定义一致,避免版本不一致导致解析错误。可通过命令查询:curl http://schema-registry:8081/subjects/schema.data.normalized-value/versions/latest修复数据库表结构
若前期因字段不匹配导致表结构缺失,可手动删除现有表,让auto.create重新根据正确Schema建表;或确保auto.evolve配置正常,让连接器自动补充缺失字段。
内容的提问来源于stack exchange,提问作者Rudonja
相关产品推荐
相关产品推荐

