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

使用Confluent BigQuery Sink连接器时新表无法写入记录的问题

问题:Kafka Connect BigQuery Sink无法为新增主题自动创建表并写入数据

环境背景

  • 本地Docker与AWS MSK环境中运行的Kafka Connect表现完全一致

当前连接器配置

{
  "name": "connector-name",
  "config": {
    "connector.class": "com.wepay.kafka.connect.bigquery.BigQuerySinkConnector",
    "allowNewBigQueryFields": "true",
    "allowBigQueryRequiredFieldRelaxation": "true",
    "autoCreateTables": "true",
    "tasks.max": "2",
    "topics": "existing-topic-name, new-topic-name",
    "project": "gcp-project-name",
    "defaultDataset": "dataset_name",
    "keyfile": "/etc/kafka/secrets/gcp-cred.key",
    "schemaRetriever": "com.wepay.kafka.connect.bigquery.retrieve.IdentitySchemaRetriever",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "<my schema registry URL goes here>",
    "value.converter.enhanced.avro.schema.support": "true",
    "schema.registry.url": "<my schema registry URL goes here>",
    "offsets.retention.minutes": 2880,
    "auto.offset.reset": "earliest",
    "enable.auto.commit": "true",
    "isolation.level": "read_committed"
  }
}

历史配置对比

6-12个月前仅针对existing-topic-name启动连接器,当时配置仅包含:

"autoCreateTables": "true"

完全没有allowNewBigQueryFields和allowBigQueryRequiredFieldRelaxation这两个属性,且连接器能自动创建对应表。

问题现象

  • 新增new-topic-name到连接器配置后,BigQuery中未自动创建对应新表,新主题的记录也无法写入
  • 6-12个月前由连接器创建的原有表,仍能正常接收旧主题的新记录
  • 已通过其他团队的命令行消费者确认,new-topic-name主题中存在有效记录

尝试过的配置组合

曾尝试以下属性的多种组合,但均未解决问题:

"allowNewBigQueryFields": "true",
"allowBigQueryRequiredFieldRelaxation": "true",
"autoCreateTables": "true",
"allowSchemaUnionization": "true"

日志信息

  • 无明确报错时,仅能看到偏移量相关INFO日志:
[2024-07-19 23:19:58,500] INFO [connector-name-here|task-0] [Consumer clientId=connector-consumer-connector-name-here-0, groupId=group id here] Setting offset for partition topic-names-here-0 to the committed offset FetchPosition{offset=14013, offsetEpoch=Optional.empty, currentLeader=LeaderAndEpoch{leader=Optional[cluster-url-here:9098 (id: 1 rack: use1-az6)], epoch=36}} (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator:820)

或类似“未找到已提交的偏移量”的提示。

  • 偶尔会出现GCP认证超时错误:
Caused by: com.google.cloud.bigquery.BigQueryException: Error getting access token for service account: connect timed out

怀疑方向

连接器配置未正确适配新增主题,新主题的记录被静默丢弃,未抛出明确错误导致无法定位问题。

内容的提问来源于stack exchange,提问作者Sultan of Swing

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 22:31:09