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

Kafka Connect JDBC Sink配置:无Schema/ID实现JSON写入PostgreSQL

Kafka Connect JDBC Sink 处理纯JSON Payload 配置问题

我需要向Kafka Topic发送无Schema结构、无Schema ID的纯JSON payload,通过JDBC Sink连接器写入PostgreSQL。目前仅能让连接器在以下两种场景下正常工作:

当前可行的两种场景

  1. JSON payload携带内嵌Schema定义(无需Schema Registry)
from datetime import datetime
from confluent_kafka import Producer
import json

payload = {
    "schema": {
        "type": "struct", 
        "fields": [
            {"type": "string", "field": "value"}, 
            {"type": "int64", "field": "value_number"}, 
            {"type": "string", "field": "timestamp"}    
        ],
        "optional": False,
        "name": "postgres-sink"
    },
    "payload": {
        "value": "example",
        "value_number": 5,
        "timestamp": str(datetime.utcnow())
    }
}
send = json.dumps(payload).encode('utf-8')

producer_conf = {
    # ... Kafka连接信息
}
producer = Producer(producer_conf)
producer.produce(topic='topic', value=send)
producer.flush()
  1. JSON payload携带Schema Registry的Schema ID
import struct
from io import BytesIO
import json
from confluent_kafka import Producer
from datetime import datetime

class _ContextStringIO(BytesIO):
    def __enter__(self):
        return self

    def __exit__(self, *args):
        self.close()
        return False

def serialize(content: dict, schema_id: int):
    with _ContextStringIO() as fp:
        fp.write(struct.pack('>bI', 0, schema_id))
        fp.write(json.dumps(content).encode('utf-8'))
        return fp.getvalue()

payload = serialize({
  "value": "example-ser",
  "value_number": 5,
  "timestamp": str(datetime.utcnow())
  }, 2)

producer_conf = {
    # ... Kafka连接信息
}
producer = Producer(producer_conf)
producer.produce(topic='topic', value=payload)
producer.flush()

目标需求:处理纯JSON payload

我需要直接发送以下格式的纯JSON:

payload = json.dumps({
    "value": "example-2",
    "value_number": 5,
    "timestamp": str(datetime.utcnow())
}).encode('utf-8')

当前配置及问题

我尝试通过指定Schema Registry中的Schema ID来配置JDBC连接器,但发送纯JSON时抛出unknown magic byte错误,配置如下:

name = '<connector name>'
config = {
    # ... JDBC连接信息
    "connector.class": 'io.confluent.connect.jdbc.JdbcSinkConnector',

    # ... 其他配置(reporter、死信队列等)

    # 转换器配置
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",   
    "value.converter": "io.confluent.connect.json.JsonSchemaConverter",
    "value.converter.schema.registry.url": "<host>",
    "value.converter.ignore.default.for.nullables": "true",
    "value.converter.schemas.enable": "false",
    "value.converter.subject.name.strategy": "io.confluent.kafka.serializers.subject.TopicNameStrategy",
    "value.converter.use.schema.id": "2",
}

res = requests.put(f'http://localhost:8083/connectors/{name}/config', json=config)
print(f"returned status code: {res.status_code} (reason: {res.reason})")
res.raise_for_status()

已确认Schema ID对应的Schema存在于Schema Registry中,但配置未生效。


问题排查与解决方案

问题根源

io.confluent.connect.json.JsonSchemaConverter的设计目标是处理两种格式的消息:

  • 带Schema Registry元数据(magic byte + Schema ID)的二进制消息
  • 内嵌Schema定义的JSON消息

它无法直接解析纯JSON字符串,这就是unknown magic byte错误的原因。另外,value.converter.use.schema.id仅在Producer端序列化时生效,Sink端反序列化场景下该参数无效。

正确配置方案

要处理纯JSON payload,需替换转换器并调整配置:

  1. 替换Value转换器
    使用原生支持纯JSON解析的org.apache.kafka.connect.json.JsonConverter替代Confluent的JsonSchemaConverter,无需依赖Schema Registry的二进制格式。

  2. 完整可用配置示例

name = '<connector name>'
config = {
    # JDBC核心配置
    "connector.class": 'io.confluent.connect.jdbc.JdbcSinkConnector',
    "topics": "topic",
    "connection.url": "jdbc:postgresql://<host>:<port>/<db>",
    "connection.user": "<user>",
    "connection.password": "<password>",
    "insert.mode": "insert", # 根据需求选择insert/upsert
    "pk.fields": "<主键字段>", # 使用upsert时必须指定
    "auto.create": "false", # 是否自动建表(生产环境谨慎开启)
    "auto.evolve": "false", # 是否自动更新表结构

    # 转换器配置
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false", # 禁用内嵌Schema,解析纯JSON

    # ... 其他配置(死信队列、reporter等)
}

res = requests.put(f'http://localhost:8083/connectors/{name}/config', json=config)
print(f"returned status code: {res.status_code} (reason: {res.reason})")
res.raise_for_status()
  1. 字段映射补充
  • 确保JSON字段名与PostgreSQL表字段名完全匹配
  • 若字段名不匹配,可通过transforms重命名:
    "transforms": "renameField",
    "transforms.renameField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
    "transforms.renameField.renames": "old_field:new_field"
    

可选:结合Schema Registry做结构验证

如果需要用Schema Registry中的Schema验证纯JSON结构,可配合SchemaValidationTransform:

"transforms": "validateSchema",
"transforms.validateSchema.type": "io.confluent.connect.transforms.SchemaValidation$Value",
"transforms.validateSchema.schema.registry.url": "<Schema Registry地址>",
"transforms.validateSchema.subject.name.strategy": "io.confluent.kafka.serializers.subject.TopicNameStrategy",
"transforms.validateSchema.schema.id": "2"

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 14:20:33