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

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

问题排查与解决

  1. 修正字段命名匹配问题
    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一致。

  2. 处理自定义枚举类型
    unit_symbol和unit_multiplier是自定义枚举,默认转换器不会自动转为字符串。添加以下配置让枚举直接转为字符串值:

    value.converter.enum.as.string: true
    
  3. 验证Schema Registry中的Schema版本
    确认Schema Registry中存储的XxxMessage Schema与本地定义一致,避免版本不一致导致解析错误。可通过命令查询:

    curl http://schema-registry:8081/subjects/schema.data.normalized-value/versions/latest
    
  4. 修复数据库表结构
    若前期因字段不匹配导致表结构缺失,可手动删除现有表,让auto.create重新根据正确Schema建表;或确保auto.evolve配置正常,让连接器自动补充缺失字段。


内容的提问来源于stack exchange,提问作者Rudonja

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 21:20:28