如何配置Kafka Sink JDBC Connector提取指定字段并忽略冗余数据
Got it, let's tackle this problem. To pull only the data and headers fields (and ignore beforeData) for your Kafka Sink target table, you'll need to chain together a few built-in Confluent Connect transforms. Here's exactly how to configure it, with breakdowns of each step:
Original CDC Message Structure
First, let's recap the incoming message we're working with:
{ "data":{ "USER_ID":{ "string":"1" }, "USER_CATEGORY":{ "string":"A" } }, "beforeData":{ "Data":{ "USER_ID":{ "string":"1" }, "USER_CATEGORY":{ "string":"B" } } }, "headers":{ "operation":"UPDATE", "timestamp":"2018-05-03T13:53:43.000" } }
Target Sink Table Fields
We need to end up with these columns in our target table:USER_ID, USER_CATEGORY, operation, timestamp
Complete Sink Connector Configuration
Add these transform settings to your Kafka Sink connector config file:
# Enable the transform chain transforms=flattenStruct,cleanUnwantedFields,renameToTargetColumns # Step 1: Flatten nested JSON structures (uses underscore as separator) transforms.flattenStruct.type=org.apache.kafka.connect.transforms.Flatten$Value transforms.flattenStruct.delimiter=_ # Step 2: Remove all fields related to beforeData (blacklist pattern) transforms.cleanUnwantedFields.type=org.apache.kafka.connect.transforms.ReplaceField$Value transforms.cleanUnwantedFields.blacklist=beforeData_* # Step 3: Rename flattened fields to match target table column names transforms.renameToTargetColumns.type=org.apache.kafka.connect.transforms.ReplaceField$Value transforms.renameToTargetColumns.renames= \ data_USER_ID_string:USER_ID, \ data_USER_CATEGORY_string:USER_CATEGORY, \ headers_operation:operation, \ headers_timestamp:timestamp
Breakdown of Each Transform
Let's walk through what each step does to the message:
- Flatten$Value: This takes the nested JSON and turns it into a flat structure. After this step, your message looks like this:
{ "data_USER_ID_string": "1", "data_USER_CATEGORY_string": "A", "beforeData_Data_USER_ID_string": "1", "beforeData_Data_USER_CATEGORY_string": "B", "headers_operation": "UPDATE", "headers_timestamp": "2018-05-03T13:53:43.000" } - ReplaceField$Value (cleanUnwantedFields): The
blacklistoption uses a wildcard to remove every field starting withbeforeData_, wiping out the old data we don't need entirely. - ReplaceField$Value (renameToTargetColumns): The
renamesmapping converts the flattened, verbose field names to the exact column names your target table expects. This ensures the sink writes data to the correct fields.
This setup works reliably assuming your CDC tool consistently outputs the {string: "value"} structure for the data fields. If your schema changes slightly, you can adjust the rename mappings accordingly.
内容的提问来源于stack exchange,提问作者Giorgos Myrianthous

