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

Kafka Connect JDBC同步MySQL时如何配置自定义Schema Registry?

Solution for Mapping MySQL JSON VARCHAR to Custom Avro Schema in Kafka Connect

Great question! You absolutely can configure your JDBC Source Connector to map that JSON-stored attributes field to your custom Avro schema without generating new schema versions. Here's a step-by-step solution:


Step 1: Pre-Register Your Custom Schema in Schema Registry

First, you need to register your target Avro schema under the existing mysql-stations-value subject (since you don't want new schema versions).

Note: If you already have an auto-generated version 1 schema for this subject, you'll need to delete it first (Schema Registry doesn't allow modifying existing versions):

curl -X DELETE http://localhost:8081/subjects/mysql-stations-value

Then register your custom schema:

curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
--data '{"schema": "{\"type\":\"record\",\"name\":\"stations\",\"namespace\":\"com.mycorp.mynamespace\",\"fields\":[{\"name\":\"code\",\"type\":\"string\"},{\"name\":\"date_measuring\",\"type\":{\"connect.name\":\"org.apache.kafka.connect.data.Timestamp\",\"connect.version\":1,\"logicalType\":\"timestamp-millis\",\"type\":\"long\"}},{\"name\":\"attributes\",\"type\":{\"type\":\"record\",\"name\":\"AttributesRecord\",\"fields\":[{\"name\":\"H1\",\"type\":\"long\",\"default\":0},{\"name\":\"H2\",\"type\":\"long\",\"default\":0},{\"name\":\"H3\",\"type\":\"long\",\"default\":0},{\"name\":\"H\",\"type\":\"long\",\"default\":0},{\"name\":\"Q\",\"type\":\"long\",\"default\":0},{\"name\":\"P1\",\"type\":\"long\",\"default\":0},{\"name\":\"P2\",\"type\":\"long\",\"default\":0},{\"name\":\"P3\",\"type\":\"long\",\"default\":0},{\"name\":\"P\",\"type\":\"long\",\"default\":0},{\"name\":\"T\",\"type\":\"long\",\"default\":0},{\"name\":\"Hr\",\"type\":\"long\",\"default\":0},{\"name\":\"pH\",\"type\":\"long\",\"default\":0},{\"name\":\"RX\",\"type\":\"long\",\"default\":0},{\"name\":\"Ta\",\"type\":\"long\",\"default\":0},{\"name\":\"C\",\"type\":\"long\",\"default\":0},{\"name\":\"OD\",\"type\":\"long\",\"default\":0},{\"name\":\"TU\",\"type\":\"long\",\"default\":0},{\"name\":\"MO\",\"type\":\"long\",\"default\":0},{\"name\":\"AM\",\"type\":\"long\",\"default\":0},{\"name\":\"N03\",\"type\":\"long\",\"default\":0},{\"name\":\"P04\",\"type\":\"long\",\"default\":0},{\"name\":\"SS\",\"type\":\"long\",\"default\":0},{\"name\":\"PT\",\"type\":\"long\",\"default\":0}]}]}"}' \
http://localhost:8081/subjects/mysql-stations-value/versions

Step 2: Update Your Connector Configuration

Modify your JDBC Source Connector config to:

  1. Disable auto-schema registration (to avoid new versions)
  2. Use a transform to parse the attributes string into a structured format that matches your custom schema
  3. Force use of the pre-registered schema

Here's the updated config:

{
  "name": "jdbc_source_mysql_stations",
  "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
  "key.converter": "io.confluent.connect.avro.AvroConverter",
  "key.converter.schema.registry.url": "http://localhost:8081",
  "value.converter": "io.confluent.connect.avro.AvroConverter",
  "value.converter.schema.registry.url": "http://localhost:8081",
  "value.converter.auto.register.schemas": "false",
  "value.converter.use.latest.version": "true",
  "transforms": ["ValueToKey", "ParseAttributes"],
  "transforms.ValueToKey.type": "org.apache.kafka.connect.transforms.ValueToKey",
  "transforms.ValueToKey.fields": ["code", "date_measuring"],
  "transforms.ParseAttributes.type": "org.apache.kafka.connect.transforms.Json$Value",
  "transforms.ParseAttributes.field": "attributes",
  "transforms.ParseAttributes.schema": "{\"type\":\"struct\",\"fields\":[{\"name\":\"H1\",\"type\":\"int64\",\"default\":0},{\"name\":\"H2\",\"type\":\"int64\",\"default\":0},{\"name\":\"H3\",\"type\":\"int64\",\"default\":0},{\"name\":\"H\",\"type\":\"int64\",\"default\":0},{\"name\":\"Q\",\"type\":\"int64\",\"default\":0},{\"name\":\"P1\",\"type\":\"int64\",\"default\":0},{\"name\":\"P2\",\"type\":\"int64\",\"default\":0},{\"name\":\"P3\",\"type\":\"int64\",\"default\":0},{\"name\":\"P\",\"type\":\"int64\",\"default\":0},{\"name\":\"T\",\"type\":\"int64\",\"default\":0},{\"name\":\"Hr\",\"type\":\"int64\",\"default\":0},{\"name\":\"pH\",\"type\":\"int64\",\"default\":0},{\"name\":\"RX\",\"type\":\"int64\",\"default\":0},{\"name\":\"Ta\",\"type\":\"int64\",\"default\":0},{\"name\":\"C\",\"type\":\"int64\",\"default\":0},{\"name\":\"OD\",\"type\":\"int64\",\"default\":0},{\"name\":\"TU\",\"type\":\"int64\",\"default\":0},{\"name\":\"MO\",\"type\":\"int64\",\"default\":0},{\"name\":\"AM\",\"type\":\"int64\",\"default\":0},{\"name\":\"N03\",\"type\":\"int64\",\"default\":0},{\"name\":\"P04\",\"type\":\"int64\",\"default\":0},{\"name\":\"SS\",\"type\":\"int64\",\"default\":0},{\"name\":\"PT\",\"type\":\"int64\",\"default\":0}]}",
  "connection.url": "jdbc:mysql://localhost:3306/db_name?useJDBCCompliantTimezoneShift=true&useLegacyDatetimeCode=false&serverTimezone=UTC",
  "connection.user": "confluent",
  "connection.password": "**************",
  "table.whitelist": ["stations"],
  "mode": "timestamp",
  "timestamp.column.name": ["date_measuring"],
  "validate.non.null": "false",
  "topic.prefix": "mysql-"
}

Key Config Explainers:

  • value.converter.auto.register.schemas: false: Prevents the connector from auto-generating new schema versions
  • value.converter.use.latest.version: true: Forces use of the pre-registered schema in Schema Registry
  • ParseAttributes transform: Converts the attributes VARCHAR string into a Kafka Connect Struct that matches your custom Avro schema's AttributesRecord structure

Fallback Option: Kafka Streams (If Transforms Don't Work)

If the above transform approach hits edge cases (e.g., complex JSON parsing logic), you can use Kafka Streams to reprocess the data:

  1. Consume the mysql-stations topic using the original auto-generated schema
  2. Parse the attributes string into your custom AttributesRecord object
  3. Serialize the full record using your pre-registered Avro schema
  4. Write the processed data back to the same topic (or a new one, if preferred)

This is more heavyweight but gives you full control over parsing logic.


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 20:52:57