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

如何配置Kafka Sink JDBC Connector提取指定字段并忽略冗余数据

Extracting Specific CDC Fields for Kafka Sink Connect

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:

  1. 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"
    }
    
  2. ReplaceField$Value (cleanUnwantedFields): The blacklist option uses a wildcard to remove every field starting with beforeData_, wiping out the old data we don't need entirely.
  3. ReplaceField$Value (renameToTargetColumns): The renames mapping 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:06:54