Kafka Connect JDBC同步MySQL时如何配置自定义Schema Registry?
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:
- Disable auto-schema registration (to avoid new versions)
- Use a transform to parse the
attributesstring into a structured format that matches your custom schema - 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 versionsvalue.converter.use.latest.version: true: Forces use of the pre-registered schema in Schema RegistryParseAttributestransform: Converts theattributesVARCHAR string into a Kafka Connect Struct that matches your custom Avro schema'sAttributesRecordstructure
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:
- Consume the
mysql-stationstopic using the original auto-generated schema - Parse the
attributesstring into your customAttributesRecordobject - Serialize the full record using your pre-registered Avro schema
- 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

