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

如何通过Tabular Iceberg Sink Connector将Kafka数据写入指定Iceberg命名空间

问题:CDC数据流同步至Iceberg时无法指定目标命名空间

架构概述

当前构建的CDC数据流链路:

  • Debezium MySQL Source Connector:捕获MySQL表的数据变更
  • Redpanda:存储变更事件消息
  • Tabular Iceberg Sink Connector:将消息同步至Nessie管理的Iceberg Catalog
  • Minio:作为Iceberg的底层存储

Kafka数据样例

{"after": {"realtimeAnalytics.inventory.customers.Value": {"balance": 441.63,"car_make": "Ford","car_year": 2013,"id": "0010130d-3130-4c96-adbb-106d3d5c0c99","owner_address": "1609 Harper Stravenue\nMoranview, SD 86752","owner_name": "Cody Graham","owner_phone_number": "174.418.3721 ext. 272","plate_number": "5961-PPP","subscription_end": "2025-04-03T00:00:00Z","subscription_start": "2024-04-03T00:00:00Z","subscription_status": "expired","timestamp": "2025-01-02T04:52:01Z"}},"before": {"realtimeAnalytics.inventory.customers.Value": {"balance": 441.63,"car_make": "Ford","car_year": 2013,"id": "0010130d-3130-4c96-adbb-106d3d5c0c99","owner_address": "1609 Harper Stravenue\nMoranview, SD 86752","owner_name": "Cody Graham","owner_phone_number": "174.418.3721 ext. 274","plate_number": "5961-PPP","subscription_end": "2025-04-03T00:00:00Z","subscription_start": "2024-04-03T00:00:00Z","subscription_status": "expired","timestamp": "2025-01-02T04:44:58Z"}},"op": "u","source": {"connector": "mysql","db": "inventory","file": "mysql-bin.000371","gtid": null,"name": "realtimeAnalytics","pos": 1573,"query": null,"row": 0,"sequence": null,"server_id": 223344,"snapshot": "false","table": "customers","thread": 14,"ts_ms": 1735793521000,"version": "2.4.2.Final"},"transaction": null,"ts_ms": 1735793521542}

Debezium MySQL Source Connector配置

"config": {"connect.keep.alive": "false","connector.class": "io.debezium.connector.mysql.MySqlConnector","database.allowPublicKeyRetrieval": "true","database.hostname": "mysql","database.include.list": "inventory","database.password": "dbz","database.port": "3306","database.server.id": "184054","database.user": "debezium","decimal.handling.mode": "double","enable.time.adjuster": "false","exactly.once.source.support": "enabled","gtid.source.filter.dml.events": "false","include.query": "false","include.schema.changes": "true","schema.history.internal.kafka.bootstrap.servers": "redpanda:9092","schema.history.internal.kafka.topic": "schema-changes.realtimeAnalytics","skipped.operations": "none","snapshot.mode": "when_needed","table.ignore.builtin": "false","time.precision.mode": "connect","tasks.max": "1","tombstones.on.delete": "true","topic.creation.default.partitions": "1","topic.creation.default.replication.factor": "-1","topic.creation.enable": "true","topic.prefix": "realtimeAnalytics","key.converter": "io.confluent.connect.avro.AvroConverter","key.converter.schema.registry.url": "http://redpanda:8081","value.converter": "io.confluent.connect.avro.AvroConverter","value.converter.schema.registry.url": "http://redpanda:8081","exactly.once.support": "requested"}

Tabular Iceberg Sink Connector当前配置

"config": {"connector.class": "io.tabular.iceberg.connect.IcebergSinkConnector","errors.logs.enabled": "true","errors.logs.include.messages": "true","topics.regex": "realtimeAnalytics.(.*)","iceberg.tables.dynamic-enabled": "true","iceberg.tables.route-field": "source.table","iceberg.tables.cdc_field": "op","iceberg.tables.auto-create-enabled": "true","iceberg.control.commit.interval-ms": "60000","iceberg.tables.auto-create-props.gc.enabled": "false","iceberg.tables.auto-create-props.write.metadata.delete-after-commit.enabled": "false","iceberg.tables.auto-create-props.format-version": "2","iceberg.tables.auto-create-props.write.delete.mode": "copy-on-write","iceberg.tables.auto-create-props.write.update.mode": "copy-on-write","iceberg.tables.auto-create-props.write.merge.mode": "copy-on-write","iceberg.tables.auto-create-props.compatibility.snapshot-id-inheritance.enabled": "true","iceberg.tables.auto-create-props.format": "PARQUET","iceberg.tables.auto-create-props.write.parquet.page-row-limit": "20000","iceberg.tables.auto-create-props.write.parquet.compression-codec": "zstd","iceberg.tables.auto-create-props.max-concurrent-file-group-rewrites": "20","iceberg.tables.auto-create-props.history.expire.max-snapshot-age-ms": "8640000000","iceberg.tables.auto-create-props.write.target-file-size-bytes": "536870912","iceberg.tables.auto-create-props.table_type": "ICEBERG","iceberg.tables.auto-create-props.partial-progress.enabled": "true","iceberg.tables.auto-create-props.write.wap.enabled": "true","iceberg.tables.evolve-schema-enabled": "true","iceberg.catalog.authentication.type": "NONE","iceberg.catalog.catalog-impl": "org.apache.iceberg.nessie.NessieCatalog","iceberg.catalog.client.region": "us-east-1","iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO","iceberg.catalog.ref": "main","iceberg.catalog.s3.access-key-id": "minio","iceberg.catalog.s3.endpoint": "http://minio:9000","iceberg.catalog.s3.path-style-access": "true","iceberg.catalog.s3.secret-access-key": "minio123","iceberg.catalog.uri": "http://nessie:19120/api/v1","iceberg.catalog.warehouse": "s3a://warehouse/NessieData","iceberg.control.commitIntervalMs": "60000","key.converter": "io.confluent.connect.avro.AvroConverter","key.converter.schema.registry.url": "http://redpanda:8081","tasks.max": "2","value.converter": "io.confluent.connect.avro.AvroConverter","value.converter.schema.registry.url": "http://redpanda:8081","value.converter.schemas.enable": "false"}

问题详情

已手动创建Nessie命名空间changelogs.inventory(对应MySQL数据库名inventory),希望将Kafka中的变更数据写入该命名空间并自动创建对应表,但尝试配置table.namespace、iceberg.catalog.default-namespace等参数后,数据始终写入Nessie Catalog的根目录,所有组件部署在本地Docker环境。

解决方案

1. 固定指定目标命名空间

在Iceberg Sink Connector的config中添加以下配置项,强制所有动态创建的表归属到changelogs.inventory命名空间:

"iceberg.tables.namespace": "changelogs.inventory"

2. 动态生成命名空间(适配多库场景)

如果后续需要同步多个MySQL数据库,可通过Kafka消息中的source.db字段动态生成命名空间,配置如下:

"iceberg.tables.route-field": "source.db",
"iceberg.tables.namespace-template": "changelogs.${source.db}",
"iceberg.tables.name-template": "${source.table}"

该配置会自动根据MySQL数据库名生成changelogs.<db_name>命名空间,并在其中创建对应表。

3. 配置优先级与权限检查

  • 确认使用的是Iceberg Sink Connector的官方配置前缀:iceberg.tables.namespace(而非旧版的table.namespace)
  • 验证Nessie中changelogs.inventory命名空间已正确创建,且Connector具备写入该命名空间的权限

4. 验证生效

修改配置后重启Sink Connector,检查Nessie Catalog中的表是否出现在changelogs.inventory命名空间下,同时查看Connector日志排查潜在报错。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 05:24:55