All Products
Search
Document Center

Realtime Compute for Apache Flink:Kafka YAML connector

Last Updated:Sep 15, 2026

The Kafka connector can be used in Flink CDC data ingestion jobs as a source or a sink. This topic describes the syntax, parameters, and examples of the Kafka YAML connector.

Prerequisites

Connect to your cluster by using one of the following methods:

  • Connect to an ApsaraMQ for Kafka cluster

    • The Kafka cluster runs version 0.11 or later.

    • ApsaraMQ for Kafka cluster is created. For more information, see Create resources.

    • The Flink workspace and the Kafka cluster reside in the same VPC, and ApsaraMQ for Kafka has added Flink to its whitelist. For more information, see Configure a whitelist.

    Important

    Limits on writing to ApsaraMQ for Kafka:

    • ApsaraMQ for Kafka does not support writing data in the zstd compression format.

    • ApsaraMQ for Kafka does not support idempotent writes or transactional writes, so the exactly-once semantics provided by Kafka sink tables are unavailable. Starting from VVR 8.0.0, the open source Kafka Client used by the Kafka Connector is upgraded to version 3.x, in which the properties.enable.idempotence property defaults to true, which explicitly enables idempotent writes. Therefore, when you write to ApsaraMQ for Kafka in VVR 8.0.0 or later, you must explicitly add properties.enable.idempotence=false to the sink table to disable idempotent writes and avoid write failures. For the storage engine comparison and feature limits of ApsaraMQ for Kafka, see Storage engine comparison.

  • Connect to a self-managed Apache Kafka cluster

    • The self-managed Apache Kafka cluster runs version 0.11 or later.

    • The network between Flink and the self-managed Apache Kafka cluster is connected. To connect to a self-managed cluster over the Internet, see Network connections.

    • Only the client configuration items of Apache Kafka 2.8 are supported. For more information, see the Apache Kafka consumer and producer configuration documentation.

Limits

  • We recommend that you use Kafka as a data source for Flink CDC data ingestion in VVR 11.1 or later.

  • Only the JSON, Debezium JSON, and Canal JSON formats are supported. Other data formats are not supported.

  • For sources, data of the same table distributed across multiple partitions is supported only in VVR 8.0.11 and later.

Usage notes

Currently, transactional writes are not recommended because of design limitations in the Flink and Kafka communities. When you set sink.delivery-guarantee = exactly-once, the Kafka Connector enables transactional writes. The following three issues are known:

  • Each Checkpoint generates a transaction ID. If the Checkpoint interval is too short, too many transaction IDs are generated. The Coordinator of the Kafka cluster may run out of memory, which compromises the stability of the Kafka cluster.

  • Each transaction creates a Producer instance. If too many transactions are committed at the same time, the TaskManager may run out of memory, which compromises the stability of the Flink job.

  • If multiple Flink jobs use the same sink.transactional-id-prefix, the transaction IDs that they generate may conflict. When one job fails to write data, the LSO (Log Start Offset) of the Kafka partition is blocked from advancing, which affects all consumers that read data from the partition.

If you require exactly-once semantics, use the Upsert Kafka connector to write data to a primary key table and rely on the primary key to ensure idempotence. If you must use transactional writes, see Exactly-once semantics.

Syntax

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: ${kafka.topic}
sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: localhost:9092

Configuration items

  • General

    Parameter

    Description

    Required

    Data type

    Default value

    Notes

    type

    The type of the source or sink.

    Yes

    String

    None

    The value must be kafka.

    name

    The name of the source or sink.

    No

    String

    None

    None

    properties.bootstrap.servers

    The address of the Kafka broker.

    Yes

    String

    None

    The format is host:port,host:port,host:port. Use commas (,) to separate multiple addresses.

    properties.*

    The configurations that are passed directly to the Kafka client.

    No

    String

    None

    The suffix must be a producer or consumer configuration item defined in the official Apache Kafka documentation.

    Flink removes the properties. prefix and passes the remaining configuration to the Kafka client. For example, you can use 'properties.allow.auto.create.topics' = 'false' to disable automatic topic creation.

    key.format

    The format used to read or write the key part of Kafka messages.

    No

    String

    None

    • For sources, only json is supported.

    • For sinks, valid values:

      • csv

      • json

    Note

    This parameter is supported only in VVR 11.0.0 and later.

    value.format

    The format used to read or write the value part of Kafka messages.

    No

    String

    debezium-json

    • For sources, valid values:

      • debezium-json 

      • canal-json

      • json

    • For sinks, valid values:

      • debezium-json 

      • canal-json

      • canal-protobuf

    Note
    • The debezium-json and canal-json formats are supported only in VVR 8.0.10 and later.

    • The json format is supported only in VVR 11.0.0 and later.

  • Source

    Parameter

    Description

    Required

    Data type

    Default value

    Notes

    topic

    The name of the topic to read.

    No

    String

    None

    Separate multiple topic names with semicolons (;), for example, topic-1 and topic-2.

    Note

    You can specify only one of topic and topic-pattern.

    topic-pattern

    The regular expression used to match the names of the topics to read. All topics that match the regular expression are read when the job runs.

    No

    String

    None

    Examples:

    • user_event_.*: matches all topics whose names start with user_event_

    • prod\.logs\..*: matches topics whose names have the prod.logs. prefix (. must be escaped)

    Note

    You can specify only one of topic and topic-pattern.

    properties.group.id

    The ID of the consumer group.

    No

    String

    None

    If the specified group ID is used for the first time, you must set properties.auto.offset.reset to earliest or latest to specify the initial startup offset.

    scan.startup.mode

    The startup offset from which Kafka reads data.

    No

    String

    group-offsets

    Valid values:

    • earliest-offset: reads data from the earliest offset of the Kafka partition.

    • latest-offset: reads data from the latest offset of the Kafka partition.

    • group-offsets (default): reads data from the committed offset of the consumer group specified by properties.group.id.

    • timestamp: reads data from the timestamp specified by scan.startup.timestamp-millis.

    • specific-offsets: reads data from the offset specified by scan.startup.specific-offsets.

    Note

    This parameter takes effect only when the job starts without state. When the job restarts from a Checkpoint or recovers from state, reading resumes from the progress saved in the state.

    scan.startup.specific-offsets

    The startup offset of each partition when the startup mode is specific-offsets.

    No

    String

    None

    Example: partition:0,offset:42;partition:1,offset:300

    scan.startup.timestamp-millis

    The startup timestamp when the startup mode is timestamp.

    No

    Long

    None

    Unit: milliseconds.

    scan.topic-partition-discovery.interval

    The interval at which Kafka topics and partitions are dynamically detected.

    No

    Duration

    5 minutes

    The default partition discovery interval is 5 minutes. To disable this feature, explicitly set the partition discovery interval to a non-positive value. When dynamic partition discovery is enabled, the Kafka source automatically discovers new partitions and reads data from them. In topic-pattern mode, the source reads data from new partitions of existing topics and from all partitions of new topics that match the regular expression.

    scan.check.duplicated.group.id

    Specifies whether to check whether the consumer group specified by properties.group.id is duplicated.

    No

    Boolean

    false

    Valid values:

    • true: checks for consumer group conflicts before the job starts. If a conflict exists, the job reports an error. This avoids conflicts with existing consumer groups.

    • false: starts the job directly without checking for consumer group conflicts.

    schema.inference.strategy

    The Schema parsing policy.

    No

    String

    continuous

    Valid values:

    • continuous: parses the Schema of every record. When two consecutive Schemas are incompatible, a broader Schema is parsed and a Schema change event is generated.

    • static: parses the Schema only once when the job starts. Data is then parsed based on the initial Schema, and no Schema change events are generated.

    Note

    scan.max.pre.fetch.records

    The maximum number of messages to consume and parse from each partition during the initial Schema parsing.

    No

    Int

    50

    Before the job reads and processes data, the specified number of latest messages is consumed in advance from each partition to initialize the Schema information.

    key.fields-prefix

    The custom prefix added to the field names parsed from the message key. This avoids naming conflicts after the Kafka message key is parsed.

    No

    String

    None

    For example, if this parameter is set to key_ and the key contains a field named a, the field is named key_a after the key is parsed.

    Note

    The value of key.fields-prefix cannot be a prefix of value.fields-prefix.

    value.fields-prefix

    The custom prefix added to the field names parsed from the message value. This avoids naming conflicts after the Kafka message value is parsed.

    No

    String

    None

    For example, if this parameter is set to value_ and the value contains a field named b, the field is named value_b after the value is parsed.

    Note

    The value of value.fields-prefix cannot be a prefix of key.fields-prefix.

    metadata.list

    The metadata columns to pass to downstream.

    No

    String

    None

    The available metadata columns are topic, partition, offset, timestamp, timestamp-type, headers, leader-epoch, __raw_key__, and __raw_value__. Separate multiple columns with commas (,).

    Note

    __raw_key__ and __raw_value__ metadata columns are available only in VVR 11.6 and later.

    scan.value.initial-schemas.ddls

    The initial Schema of specific tables, specified by using DDL statements.

    No

    String

    None

    Separate multiple DDL statements with ;. For example, use CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT); to specify the initial Schemas of the db1.t1 and db1.t2 tables.

    The table Schema in the DDL statement must be consistent with the destination table and comply with Flink SQL syntax rules.

    Note

    This parameter is supported only in VVR 11.5 and later.

    ingestion.ignore-errors

    Specifies whether to ignore errors during data parsing.

    No

    Boolean

    false

    Note

    This parameter is supported only in VVR 11.5 and later.

    ingestion.error-tolerance.max-count

    The number of accumulated parsing errors after which the job fails when data parsing errors are ignored.

    No

    Integer

    -1

    This parameter takes effect only when ingestion.ignore-errors is enabled. The default value -1 indicates that parsing exceptions do not cause the job to fail.

    Note

    This parameter is supported only in VVR 11.5 and later.

    scan.duplicate-field.strategy

    Specifies how to handle duplicate field names parsed from the key and value parts.

    No

    String

    EXCEPTION

    Valid values:

    • EXCEPTION: throws an exception when duplicate fields exist in the key and value. This is the default behavior in VVR 11.6 and earlier.

    • PREFER_KEY: prefers the value of the key field when fields are duplicated.

    • PREFER_VALUE: prefers the value of the value field when fields are duplicated.

    Note

    This parameter is supported only in VVR 11.7 and later.

    • Source table in Debezium JSON format

      Parameter

      Required

      Data type

      Default value

      Description

      debezium-json.distributed-tables

      No

      Boolean

      false

      Enable this option if the data of a single table in the Debezium JSON appears in multiple partitions.

      Note

      This parameter is supported only in VVR 8.0.11 and later.

      Important

      After you change this parameter, you must start the job without state.

      debezium-json.schema-include

      No

      Boolean

      false

      When you configure Debezium Kafka Connect, you can enable the Kafka configuration value.converter.schemas.enable to include the schema in messages. This option specifies whether Debezium JSON messages contain the schema.

      Valid values:

      • true: Debezium JSON messages contain the schema.

      • false: Debezium JSON messages do not contain the schema.

      debezium-json.ignore-parse-errors

      No

      Boolean

      false

      Valid values:

      • true: skips the current row when a parsing exception occurs.

      • false (default): reports an error, and the job fails to start.

      debezium-json.infer-schema.primitive-as-string

      No

      Boolean

      false

      Specifies whether to parse all types as String when the table Schema is parsed.

      Valid values:

      • true: parses all primitive types as String.

      • false (default): parses data based on the basic rules.

      debezium-json.infer-schema.string-type-inference

      No

      Boolean

      true

      Specifies whether to attempt to infer string fields as the TIME, DATE, or TIMESTAMP type. If this parameter is set to false, the inference is skipped and the fields remain STRING.

      Note

      This parameter is supported only in VVR 11.8 and later.

    • Source table in Canal JSON format

      Parameter

      Required

      Data type

      Default value

      Description

      canal-json.distributed-tables

      No

      Boolean

      false

      Enable this option if the data of a single table in the Canal JSON appears in multiple partitions.

      Note

      This parameter is supported only in VVR 8.0.11 and later.

      Important

      After you change this parameter, you must start the job without state.

      canal-json.database.include

      No

      String

      None

      An optional regular expression that matches the database metadata field in Canal records so that only the changelog records of the specified database are read. The regular expression string is compatible with Java Pattern.

      canal-json.table.include

      No

      String

      None

      An optional regular expression that matches the table metadata field in Canal records so that only the changelog records of the specified table are read. The regular expression string is compatible with Java Pattern.

      canal-json.ignore-parse-errors

      No

      Boolean

      false

      Valid values:

      • true: skips the current row when a parsing exception occurs.

      • false (default): reports an error, and the job fails to start.

      canal-json.infer-schema.primitive-as-string

      No

      Boolean

      false

      Specifies whether to parse all types as String when the table Schema is parsed.

      Valid values:

      • true: parses all primitive types as String.

      • false (default): parses data based on the basic rules.

      canal-json.infer-schema.strategy

      No

      String

      AUTO

      The parsing policy used when the table Schema is parsed.

      Valid values:

      • AUTO (default): parses types automatically by parsing the JSON data. If the data does not contain the sqlType field, we recommend that you use AUTO to avoid parsing failures.

      • SQL_TYPE: parses types based on the sqlType array in the Canal JSON data. If the data contains the sqlType field, we recommend that you set canal-json.infer-schema.strategy to SQL_TYPE for more precise types.

      • MYSQL_TYPE: parses types based on the mysqlType array in the Canal JSON data.

      When the Canal JSON data in Kafka contains the sqlType field and you require more precise type mapping, we recommend that you set canal-json.infer-schema.strategy to SQL_TYPE.

      For the sqlType type mapping rules, see Canal JSON schema parsing.

      Note
      • This parameter is supported only in VVR 11.1 and later.

      • MYSQL_TYPE is supported only in VVR 11.3 and later.

      canal-json.mysql.treat-mysql-timestamp-as-datetime-enabled

      No

      Boolean

      true

      Specifies whether to map the MySQL timestamp type to the CDC timestamp type:

      • true (default): maps the MySQL timestamp type to the CDC timestamp type.

      • false: maps the MySQL timestamp type to the CDC timestamp_ltz type.

      canal-json.mysql.treat-tinyint1-as-boolean.enabled

      No

      Boolean

      true

      Specifies whether to map the MySQL tinyint(1) type to the CDC boolean type when MYSQL_TYPE parsing is used:

      • true (default): maps the MySQL tinyint(1) type to the CDC boolean type.

      • false: maps the MySQL tinyint(1) type to the CDC tinyint(1) type.

      This parameter takes effect only when canal-json.infer-schema.strategy is set to MYSQL_TYPE.

      canal-json.infer-schema.string-type-inference

      No

      Boolean

      true

      Specifies whether to attempt to infer string fields as the TIME, DATE, or TIMESTAMP type. If this parameter is set to false, the inference is skipped and the fields remain STRING.

      Note

      This parameter is supported only in VVR 11.8 and later.

    • Source table in JSON format

      Parameter

      Required

      Data type

      Default value

      Description

      json.timestamp-format.standard

      No

      String

      SQL

      The input and output timestamp format. Valid values:

      • SQL: parses input timestamps in the yyyy-MM-dd HH:mm:ss.s{precision} format, for example, 2020-12-30 12:13:14.123.

      • ISO-8601: parses input timestamps in the yyyy-MM-ddTHH:mm:ss.s{precision} format, for example, 2020-12-30T12:13:14.123.

      json.ignore-parse-errors

      No

      Boolean

      false

      Valid values:

      • true: skips the current row when a parsing exception occurs.

      • false (default): reports an error, and the job fails to start.

      json.infer-schema.primitive-as-string

      No

      Boolean

      false

      Specifies whether to parse all types as String when the table Schema is parsed.

      Valid values:

      • true: parses all primitive types as String.

      • false (default): parses data based on the basic rules.

      json.infer-schema.flatten-nested-columns.enable

      No

      Boolean

      false

      Specifies whether to recursively expand nested columns in the JSON data during parsing. Valid values:

      • true: expands nested columns recursively.

      • false (default): treats nested columns as String.

      json.decode.parser-table-id.fields

      No

      String

      None

      Specifies whether to use the values of specific JSON fields to generate the tableId when JSON data is parsed. Separate multiple fields with ,. For example, if the JSON data is {"col0":"a", "col1","b", "col2","c"}, the generated result is as follows:

      Configuration

      tableId

      col0

      a

      col0,col1

      a.b

      col0,col1,col2

      a.b.c

      json.infer-schema.fixed-types

      No

      String

      None

      The specific types of certain fields when JSON data is parsed. Separate multiple fields with ,. For example, id BIGINT, name VARCHAR(10) specifies the type of the id field in the JSON data as BIGINT and the type of the name field as VARCHAR(10).

      Note
      • This parameter is supported only in VVR 11.5 and later.

      • When you use this parameter in VVR 11.5, you must also add scan.max.pre.fetch.records: 0.

      json.decode.converter-class

      No

      String

      None

      The fully qualified name of the implementation class. The converter can modify the JSON byte[] before the JSON data is parsed.

      Note

      This parameter is supported only in VVR 11.6 and later.

      json.decode.empty-value-as-delete.enabled

      No

      Boolean

      false

      Specifies whether to parse tombstone messages (with an empty value) in Kafka compacted topics as DELETE events. This applies to scenarios in which an empty value indicates deletion, such as compacted topic mirroring and CDC deletion signals.

      Note

      This parameter is supported only in VVR 11.7 and later.

      json.infer-schema.string-type-inference

      No

      Boolean

      true

      Specifies whether to attempt to infer string fields as the TIME, DATE, or TIMESTAMP type. If this parameter is set to false, the inference is skipped and the fields remain STRING.

      Note

      This parameter is supported only in VVR 11.8 and later.

  • Sink

    Parameter

    Description

    Required

    Data type

    Default value

    Notes

    type

    The type of the sink.

    Yes

    String

    None

    The value must be kafka.

    name

    The name of the sink.

    No

    String

    None

    None

    topic

    The name of the Kafka topic.

    No

    String

    None

    If this option is enabled, all data is written to this topic.

    Note

    If this option is disabled, each record is written to the topic named after its TableID string (generated by joining with .), for example, databaseName.tableName.

    partition.strategy

    The policy used to write data to Kafka partitions.

    No

    String

    all-to-zero

    Valid values:

    • all-to-zero (default): writes all data to partition 0.

    • hash-by-key: writes data to multiple partitions based on the hash value of the primary key. This ensures that data with the same primary key is written to the same partition in order.

    sink.tableId-to-topic.mapping

    The mapping between source table names and destination Kafka topic names.

    No

    String

    None

    Separate each mapping with ;. Separate the source table name and the destination Kafka topic name with :. Table names can be regular expressions. Multiple tables that are mapped to the same topic can be joined with ,. For example: mydb.mytable1:topic1;mydb.mytable2:topic2.

    Note

    This parameter allows you to change the mapped topic while preserving the original table name information.

    • Sink table in Canal JSON format

      Parameter

      Required

      Data type

      Default value

      Description

      canal-json.serialize.update.keep-changed-fields-only

      No

      Boolean

      false

      Specifies whether the old part of an UPDATE message in the canal-json format contains only the pre-change values of the fields that changed.

      Note

      This parameter is supported only in VVR 11.8 and later.

    • Sink table in Debezium JSON format

      Parameter

      Required

      Data type

      Default value

      Description

      debezium-json.include-schema.enabled

      No

      Boolean

      false

      Specifies whether the Debezium JSON data contains Schema information.

      debezium-json.emit.full-table-id.enabled

      No

      Boolean

      false

      Specifies whether to write the complete three-part Table ID into the Debezium JSON metadata fields.

      The mapping when this parameter is enabled:

      CDC Table ID part

      Debezium JSON key

      Namespace

      db

      Schema

      schema

      Table

      table

      The mapping when this parameter is disabled:

      CDC Table ID part

      Debezium JSON key

      Namespace

      None

      Schema

      db

      Table

      table

      Note

      This parameter is supported only in VVR 11.6 and later.

Reuse an existing catalog

Starting from VVR 11.5, you can directly reference the built-in Kafka Catalog that is created on the Catalogs page in a Flink CDC data ingestion job. This reduces the need to write connection properties manually.

source:
  type: kafka
  using.built-in-catalog: kafka_catalog

Currently, data ingestion jobs support the automatic reuse of the following Kafka Catalog parameters:

  • properties.bootstrap.servers

  • format

  • key.fields-prefix

  • value.fields-prefix

  • timestamp-format.standard

  • infer-schema.flatten-nested-columns.enable

  • infer-schema.primitive-as-string

  • max.fetch.records

To override the automatically reused parameters, explicitly specify the corresponding YAML parameters, which take higher priority.

Configuration examples

  • Use Kafka as the source of a data ingestion job:

    source:
      type: kafka
      name: Kafka source
      properties.bootstrap.servers: ${kafka.bootstraps.server}
      topic: ${kafka.topic}
      value.format: ${value.format}
      scan.startup.mode: ${scan.startup.mode}
     
    sink:
      type: hologres
      name: Hologres sink
      endpoint: <yourEndpoint>
      dbname: <yourDbname>
      username: ${secret_values.ak_id}
      password: ${secret_values.ak_secret}
      sink.type-normalize-strategy: BROADEN
  • Use Kafka as the sink of a data ingestion job:

    source:
      type: mysql
      name: MySQL Source
      hostname: ${secret_values.mysql.hostname}
      port: ${mysql.port}
      username: ${secret_values.mysql.username}
      password: ${secret_values.mysql.password}
      tables: ${mysql.source.table}
      server-id: 8601-8604
    
    sink:
      type: kafka
      name: Kafka Sink
      properties.bootstrap.servers: ${kafka.bootstraps.server}
    
    route:
      - source-table: ${mysql.source.table}
        sink-table: ${kafka.topic}

    In this example, the route module is used to set the name of the Kafka topic to which the source table is written.

Note

ApsaraMQ for Kafka does not enable automatic topic creation by default. For more information, see FAQ about automatic topic creation. Before you write data to ApsaraMQ for Kafka, you must create the corresponding topics. For more information, see Step 3: Create resources.

Examples

The following examples show the configurations for typical scenarios.

Read a single topic

The following example reads the topic customers and writes data to Data Lake Formation:

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true

In this example, the table name generated for the JSON format is the same as the topic name by default.

Read multiple topics

The following example reads multiple topics that match a regular expression and writes data to the StarRocks connector:

source:
  type: kafka
  topic-pattern: user_event_.*
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous

sink:
  type: starrocks
  jdbc-url: jdbc:mysql://<yourFeHostname>:9030
  load-url: <yourFeHostname>:8030
  username: <yourUsername>
  password: ${secret_values.starrocks_password}
 
  # (可选)数据量不大的作业建议调低 flush 间隔,避免数据长时间不落盘(默认 300000,即 5 分钟)
  sink.buffer-flush.interval-ms: 5000
  # (可选)上游为 utf8mb4 字符集时建议设为 4,避免文本截断(默认 3)
  unicode-char.max-bytes: 4
  # (可选)自动建表的分桶数;StarRocks 2.5.7 以下必须显式配置,更高版本可自动推断
  table.create.num-buckets: 8
  # (可选)自动建表的副本数,按集群情况配置
  table.create.properties.replication_num: 3
  # (可选)StarRocks 3.2 及以上建议开启,加速表结构变更
  table.create.properties.fast_schema_evolution: true
  # 注意:通过 transform 变更主键时,必须同时设置 sink.ignore.update-before: false,
  # 否则旧主键对应的行会残留在下游

In this example, the table name generated for the JSON format is the same as the topic name by default.

Read keys and prevent field conflicts

To prevent errors caused by field conflicts between the key and the value, use one of the following solutions:

  1. Add a prefix to fields to avoid field conflicts:

source:
  type: kafka
  topic: ${kafka.topic}
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  key.format: json
  value.format: json
  # key部分的字段名添加key_前缀
  key.fields-prefix: key_
  # value部分的字段名添加value_前缀
  value.fields-prefix: value_
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true
  1. Configure a conflict resolution policy. For more information, see scan.duplicate-field.strategy. In the following configuration, the fields in the key take precedence, and fields with the same name in the value are ignored when duplicates occur.

source:
  type: kafka
  topic: ${kafka.topic}
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  key.format: json
  value.format: json
  # 优先使用key的字段,忽略同名的value中字段
  scan.duplicate-field.strategy: PREFER_KEY
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true

Add metadata columns

The following example reads the topic customers and adds the metadata columns topic and partition to the fields:

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous
  metadata.list: topic,partition

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true

Handle parsing errors

Data parsing errors cause job failures. You can configure a job to tolerate parsing errors, which is commonly used together with Dirty data collection.

The following example completely ignores parsing errors:

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous
  # 开启忽略解析报错,默认忽略全部解析报错
  ingestion.ignore-errors: true

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true
  
# 开启脏数据收集器,打印解析失败数据
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

The following example makes the job fail after 30 accumulated parsing failures:

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous
  # 开启忽略解析报错
  ingestion.ignore-errors: true
  # 解析报错发生30次后触发作业失败
  ingestion.error-tolerance.max-count: 30
  
sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true
  
# 开启脏数据收集器,打印解析失败数据
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Read JSON data

The following sections describe common ways to read data in the JSON format.

Specify the table id for parsing

By default, the table id of JSON data is the topic name. You can specify the values of fields in the data as the table id. In the following example, the db and tbl fields are specified as the table id.

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous
  # 指定字段 db 和 table 作为 table id
  value.json.decode.parser-table-id.fields: db,tbl

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true
  
# 开启脏数据收集器,打印解析失败数据
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Specify field types

Field types are inferred from field values. Some inferred types may not be the types that you expect. You can specify fixed types for certain fields and skip the subsequent inference and evolution of these field types.

The following configuration fixes the types of four fields:

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # 指定字段 db 和 tbl 作为 table id
  value.json.decode.parser-table-id.fields: db,tbl
  # 固定指定部分字段类型
  value.json.infer-schema.fixed-types: 'db STRING, tbl STRING, id BIGINT, amount DECIMAL(18, 2)'
  # 允许未声明字段继续动态推断
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true
  
# 开启脏数据收集器,打印解析失败数据
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Read Canal JSON data

The following sections describe common ways to read data in the Canal JSON format.

Type inference policies

By default, when Canal JSON data is read, the data values are used to infer the Schema types.

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: canal-json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true
  
# 开启脏数据收集器,打印解析失败数据
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

You can also infer types based on the Schema information (sql type or mysql type) recorded in the Canal JSON data. The following example uses mysql type to infer the Schema.

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: canal-json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous
  # 使用 mysql type 信息推导 Schema,也可以配置为 SQL_TYPE 通过 sql type 推导
  value.canal-json.infer-schema.strategy: MYSQL_TYPE

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true
  
# 开启脏数据收集器,打印解析失败数据
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Read Debezium JSON data

The following example reads the topic customers and writes data to Data Lake Formation:

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: debezium-json
  # (可选)动态识别每条数据Schema,并比对生成Schema变更
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
  commit.user: your_job_name
  #(可选)开启删除向量,提升读取性能
  table.properties.deletion-vectors.enabled: true

Synchronize raw MySQL binary logs to Kafka

Flink CDC data ingestion supports synchronizing raw MySQL binary log content to Canal JSON. The following job synchronizes the binary logs of multiple tables to the topic order_dw_tables.

source:
  type: mysql
  hostname: #{hostname}
  port: 3306
  username: #{username}
  password: #{password}
  tables: order_dw.\.*
  server-id: 28601-28604
  #(可选)同步增量阶段新创建的表的数据
  scan.binlog.newly-added-table.enabled: true
  #(可选)同步表注释和字段注释
  include-comments.enabled: true
  #(可选)优先分发无界的分片以避免可能出现的TaskManager OutOfMemory问题
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  #(可选)开启解析过滤,加速读取
  scan.only.deserialize.captured.tables.changelog.enabled: true
  # 在 Canal JSON 中补充 mysqlType、sqlType、sql、isDdl 等信息
  include-binlog-meta.enable: true
  
sink:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: order_dw_tables
  # Kafka value 使用 Canal JSON changelog 格式
  value.format: canal-json
  # 指定序列化日期类型数据时使用的格式
  value.canal-json.timestamp-format.standard: SQL
  # 数据统一写入分区0,保证binlog顺序
  partition.strategy: all-to-zero

Advanced features

Schema parsing and change synchronization policies

The Kafka connector maintains the Schemas of all currently known tables.

Initialize table Schema information

Table Schema information includes field and data type information, database and table information, and primary key information. These three types of information are initialized as follows:

  • Field and data type information

A data ingestion job can automatically infer the fields and data types of a table from the data. However, in some scenarios, you may want to specify the fields and types of certain tables. Based on the granularity at which you specify field types, table Schema information can be initialized by using one of the following three policies:

  1. Fully inferred by the program

Before the Kafka data is read, the Kafka connector attempts to consume up to the number of messages specified by scan.max.pre.fetch.records from each partition in advance, parses the Schema of each record, and then merges these Schemas to initialize the table Schema information. Before the data is actually consumed, the corresponding table creation events are generated based on the initialized Schema.

Note

For the Debezium JSON and Canal JSON formats, the table information is contained in individual messages. The scan.max.pre.fetch.records messages consumed in advance may contain data of multiple tables, so the number of records consumed in advance cannot be determined for each table. Pre-consumption and table Schema initialization are performed only once before the messages of each partition are actually consumed and processed. If data of a new table arrives later, the table Schema parsed from the first record of that table is used as the initial Schema, and pre-consumption and initialization are not performed again for that table.

Important

Data of a single table distributed across multiple partitions is supported only in VVR 8.0.11 and later. In this scenario, you must set debezium-json.distributed-tables or canal-json.distributed-tables to true.

  1. Specify the initial table Schema

In some scenarios, you may want to specify the initial table Schema yourself, for example, when you write Kafka data to a pre-created downstream table. In this case, you can add the scan.value.initial-schemas.ddls parameter to specify the initial table Schema. The following example shows a configuration:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # 使用数据中的 db、table 字段作为 Table ID
  json.decode.parser-table-id.fields: db,table
  # 设置初始表结构
  scan.value.initial-schemas.ddls: CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT);

The CREATE TABLE statement must be consistent with the Schema of the destination table. In this example, the initial type of the id field in the db1.t1 table is specified as BIGINT and the initial type of the name field is specified as VARCHAR(10), and the initial type of the id field in the db1.t2 table is specified as BIGINT.

The CREATE TABLE statement uses Flink SQL syntax.

  1. Specify fixed types for fields

In some scenarios, you may want to fix the data type of specific fields. For example, you may want certain fields that could be inferred as the TIMESTAMP type to be delivered as strings. In this case, you can add the json.infer-schema.fixed-types parameter to specify the initial table Schema. This takes effect only when the message format is json. The following example shows a configuration:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # 设置特定字段始终为固定类型
  json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
  scan.max.pre.fetch.records: 0

In this example, the type of all id fields is fixed as BIGINT and the type of all name fields is fixed as VARCHAR(10).

The types used here are the same as Flink SQL data types.

  • Database and table information

    • For the Canal JSON and Debezium JSON formats, the table information is parsed from individual messages and includes the database name and the table name.

    • For the JSON format, the table information contains only the table name by default, which is the name of the topic where the data resides. If your data contains database and table information, you can use json.infer-schema.fixed-types to specify the fields that contain the database and table information. These fields are mapped to the database name and the table name. The following example shows a configuration:

      source:
        type: kafka
        name: Kafka Source
        properties.bootstrap.servers: host:9092
        topic: test-topic
        value.format: json
        scan.startup.mode: earliest-offset
        # 使用 col1 字段中的值作为库名,使用 col2 字段中的值作为表名
        json.decode.parser-table-id.fields: col1,col2

      In this example, each record is delivered to the table whose database name is the value of the col1 field and whose table name is the value of the col2 field.

  • Primary key information

    • For the Canal JSON format, the primary key of a table is defined based on the pkNames field in the JSON.

    • For the Debezium JSON and JSON formats, the JSON does not contain primary key information. You can use transform rules to manually add primary keys to a table:

      transform:
        - source-table: \.*.\.*
          projection: \*
          primary-keys: key1, key2

Schema parsing and Schema changes

After the table Schema is initialized, if schema.inference.strategy is set to static, the Kafka connector parses the value of each message based on the initial table Schema and does not generate Schema change events. If schema.inference.strategy is set to continuous, the Kafka connector parses the value of each Kafka message, obtains the physical columns of the message, and compares them with the currently maintained Schema. If the parsed Schema is inconsistent with the current Schema, the connector attempts to merge the Schemas and generates the corresponding table Schema change events. The merge rules are as follows:

  • If the parsed physical columns contain fields that do not exist in the current Schema, these fields are added to the Schema and an add-nullable-column event is generated.

  • If the parsed physical columns do not contain a field that already exists in the current Schema, the field is retained and the data of that column is filled with NULL. No drop-column event is generated.

  • If both contain columns with the same name, the columns are handled as follows:

    • When the types are the same but the precisions differ, the type with the higher precision is used and a column-type-change event is generated.

    • When the types differ, the lowest common parent node in the following tree structure is used as the type of the column, and a column-type-change event is generated.

      image

  • The following Schema change policies are currently supported:

    • Add a column: The corresponding column is added to the end of the current Schema, the data of the new column is synchronized, and the new column is set as a nullable column.

    • Drop a column: No drop-column event is generated. Instead, the data of the column is automatically filled with NULL values.

    • Rename a column: This is treated as adding a column and dropping a column. The renamed column is added to the end of the current Schema, and the data of the column before renaming is filled with NULL values.

    • Change a column type:

      • For downstream systems that support column type changes, data ingestion jobs support type changes of regular columns after the downstream sink supports handling column type changes. For example, a column can be changed from the INT type to the BIGINT type. Such changes depend on the column type change rules supported by the downstream sink. Different sink tables support different rules. Refer to the sink table documentation for the column type change rules that it supports.

      • For downstream systems that do not support column type changes, such as Hologres, you can use Type broadening. In this case, a table with broader types is created in the downstream system when the job starts. When a column type change occurs, the job checks whether the downstream sink can accept the change, which implements tolerant support for column type changes.

  • The following Schema changes are not supported:

    • Changes to constraints such as primary keys or indexes.

    • Changes from NOT NULL to NULLABLE.

  • Canal JSON schema parsing

    Canal JSON data may contain an optional sqlType field that records the precise type information of data columns. To obtain a more accurate Schema, you can set canal-json.infer-schema.strategy to SQL_TYPE to use the types in sqlType. The type mappings are as follows:

    JDBC type

    Type Code

    CDC type

    BIT

    -7

    BOOLEAN

    BOOLEAN

    16

    TINYINT

    -6

    TINYINT

    SMALLINT

    -5

    SMALLINT

    INTEGER

    4

    INT

    BIGINT

    -5

    BIGINT

    DECIMAL

    3

    DECIMAL(38,18)

    NUMERIC

    2

    REAL

    7

    FLOAT

    FLOAT

    6

    DOUBLE

    8

    DOUBLE

    BINARY

    -2

    BYTES

    VARBINARY

    -3

    LONGVARBINARY

    -4

    BLOB

    2004

    DATE

    91

    DATE

    TIME

    92

    TIME

    TIMESTAMP

    93

    TIMESTAMP

    CHAR

    1

    STRING

    VARCHAR

    12

    LONGVARCHAR

    -1

    Other types

Dirty data tolerance and collection

In some cases, your Kafka data source may contain malformed data (dirty data). To prevent the synchronization job from failing and restarting frequently because of this dirty data, you can configure the job to ignore such abnormal data. The following example shows a configuration:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # 开启脏数据容忍功能
  ingestion.ignore-errors: true
  # 容忍 1000 条脏数据
  ingestion.error-tolerance.max-count: 1000

This configuration ignores up to 1,000 dirty records, which allows the job to run normally when a small amount of dirty data exists. When the number of dirty records exceeds this threshold, the job fails, which reminds you to validate the data.

If you do not want the job to fail because of dirty data, use the following configuration:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # 开启脏数据容忍功能
  ingestion.ignore-errors: true
  # 容忍所有的脏数据
  ingestion.error-tolerance.max-count: -1

The dirty data tolerance policy prevents the job from failing frequently because of abnormal data. You may also want to learn more about the dirty data to adjust the behavior of the Kafka data producer. By following the process described in Dirty data collection, you can view the dirty data of the job in the TaskManager logs. The following example shows a configuration:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # 开启脏数据容忍功能
  ingestion.ignore-errors: true
  # 容忍所有的脏数据
  ingestion.error-tolerance.max-count: -1

pipeline:
  dirty-data.collector:
    # 将脏数据写入 TaskManager 的日志文件中
    type: logger

Mapping policies between table names and topics

When Kafka is used as the sink of a data ingestion job, the Kafka message format (debezium-json or canal-json) also contains table name information. When the Kafka messages are consumed later, the table name in the data is usually used as the actual table name instead of the topic name. Therefore, you must carefully configure the mapping policy between table names and topics.

Assume that the two tables mydb.mytable1 and mydb.mytable2 in MySQL need to be synchronized. The following configuration policies are available:

1. Do not configure any mapping policy

Without any mapping policy, each table is written to the corresponding topic named in the database name.table name format. Therefore, the data of mydb.mytable1 is written to the topic named mydb.mytable1, and the data of mydb.mytable2 is written to the topic named mydb.mytable2. The following example shows a configuration:

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}

2. Configure route rules for mapping (not recommended)

In many scenarios, you do not want the destination topic to be in the database name.table name format and want to write data to a specified topic. In this case, you can configure route rules for mapping. The following example shows a configuration:

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  
 route:
  - source-table: mydb.mytable1,mydb.mytable2
    sink-table: mytable1

In this case, all data from mydb.mytable1 and mydb.mytable2 is written to the single topic mytable1.

However, when route rules are used to change the destination topic name, the table name information in the Kafka messages (in the debezium-json or canal-json format) is also changed. In this case, all table names in the Kafka messages become mytable1, which may cause unexpected results when other systems consume the Kafka messages of this topic.

3. Configure the sink.tableId-to-topic.mapping parameter for mapping (recommended)

To configure the mapping rules between table names and topics while preserving the source table name information, you can use the sink.tableId-to-topic.mapping parameter. The following example shows a configuration:

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  sink.tableId-to-topic.mapping: mydb.mytable1,mydb.mytable2:mytable

Or

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  sink.tableId-to-topic.mapping: mydb.mytable1:mytable;mydb.mytable2:mytable

In this case, all data from mydb.mytable1 and mydb.mytable2 is written to the single topic mytable1, and the table name information in the Kafka messages (in the debezium-json or canal-json format) remains mydb.mytable1 or mydb.mytable2. When other systems consume the Kafka messages of this topic, they can correctly obtain the source table name information.

JSON Converter implementation

JSON data may not always have a fully consistent format, or some processing requirements cannot be met by the Transform module. For example, you may need to merge two columns into a new column, drop the original two columns, and ensure that Schema change synchronization still works. To process JSON data more flexibly, you can implement the KafkaPayloadConverter interface to modify the JSON data before the framework processes it. To use this feature, add the json.decode.converter-class configuration and set it to the fully qualified name of the implementation class.

The process of implementing and using a JSON Converter is as follows:

  1. Implement the KafkaPayloadConverter interface and package it. The open source demo repository already provides some implementation examples.

  2. Add the packaged file to the additional dependencies of the data ingestion job.

  3. In the Source module, specify the json.decode.converter-class configuration. The following example uses the ArrayElementExtractorConverter class from the open source project:

    source:                                                                                                                                                                                                                                                                                                            
      type: kafka                                                                                                                                                                                                                                                                                          
      properties.bootstrap.servers: localhost:9092                                                                                                                                                                                                                                                                     
      topic: my_cdc_topic                                                                                                                                                                                                                                                                                              
      properties.group.id: flink-cdc-group                                                                                                                                                                                                                                                                             
      scan.startup.mode: earliest-offset  
      value.format: json                                                                                                                                                                                                                                                                              
      # 对 Kafka 消息 value 的 JSON 数据应用自定义转换器
      value.json.decode.converter-class: org.apache.flink.cdc.connectors.kafka.ArrayElementExtractorConverter
  4. Deploy and run the job.

key format and value format

In VVR 11.5 and earlier, the configurations of key format and value format cannot be distinguished. When the same format is used, the format configuration applies to both the key format and the value format.

VVR 11.6 and later optimize this behavior. Taking the json format as an example, the format configuration items are passed as follows:

  • format prefix (such as json.infer-schema.primitive-as-string): To maintain version compatibility, configurations with the format prefix apply to both the key format and the value format by default.

  • key prefix plus format prefix (such as key.json.infer-schema.primitive-as-string): Configurations with the key prefix plus the format prefix apply only to the key format and take higher priority than configurations with the format prefix.

  • value prefix plus format prefix (such as value.json.infer-schema.primitive-as-string): Configurations with the value prefix plus the format prefix apply only to the value format and take higher priority than configurations with the format prefix.

In the following Kafka source configuration, different JSON Converters are configured for the key format and the value format.

source:                                                                                                                                                                                                                                                                                                            
  type: kafka                                                                                                                                                                                                                                                                                          
  properties.bootstrap.servers: localhost:9092                                                                                                                                                                                                                                                                     
  topic: my_cdc_topic                                                                                                                                                                                                                                                                                              
  properties.group.id: flink-cdc-group                                                                                                                                                                                                                                                                             
  scan.startup.mode: earliest-offset 
  key.format: json  
  value.format: json                                                                                                                                                                                                                                                                            
  # 对 Kafka 消息 key 的 JSON 数据应用自定义转换器
  key.json.decode.converter-class: com.example.KeyExampleConverter                                                                                                                                                                                                                                                                            
  # 对 Kafka 消息 value 的 JSON 数据应用自定义转换器
  value.json.decode.converter-class: com.example.ValueExampleConverter