All Products
Search
Document Center

Realtime Compute for Apache Flink:Real-time log ingestion into a data lake

Last Updated:Apr 29, 2026

This topic describes the best practices for writing log data to Alibaba Cloud Data Lake Formation (DLF) with a Flink CDC data ingestion YAML job.

Kafka log ingestion into a data lake

Use Flink CDC data ingestion and a simple YAML job to quickly ingest log data into a data lake in real time. The system automatically performs schema inference and supports schema evolution.

Assume that the inventory topic in Apache Kafka stores data for a log table in JSON format. The following example job synchronizes this data to the corresponding sink table in Data Lake Formation (DLF):

source:
  type: kafka
  name: Kafka Source
  # The Kafka broker addresses.
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  # The topic to consume.
  topic: inventory
  # Specifies that consumption starts from the earliest offset.
  scan.startup.mode: earliest-offset
  # The format of the Kafka message value.
  value.format: json
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # (Optional) Recursively flattens nested columns in JSON data.
  json.infer-schema.flatten-nested-columns.enable: true
  # (Optional) Skips the first 100 parsing exceptions. If the number of exceptions exceeds 100, the job fails.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  # The metastore type. The value is fixed to rest.
  catalog.properties.metastore: rest
  # The token provider. The value is fixed to dlf.
  catalog.properties.token.provider: dlf
  # The URI to access the DLF Rest Catalog Server. The format is http://[region-id]-vpc.dlf.aliyuncs.com. For example, http://cn-hangzhou-vpc.dlf.aliyuncs.com.
  catalog.properties.uri: dlf_uri
  # The name of the DLF catalog.
  catalog.properties.warehouse: your_warehouse
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

# Adds primary key information to the table. 
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
    
# Writes all data from the inventory topic to the test_database.inventory table.
route:
  - source-table: inventory
    sink-table: test_database.inventory
    
pipeline:
  # (Optional) Records dirty data that causes processing exceptions in the logs.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Assume that the inventory topic in Apache Kafka stores data for multiple log tables in JSON format. The databaseName and tableName fields in the JSON payload provide the database and table names. The following example job synchronizes data from these tables to their corresponding sink tables in DLF:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (Optional) Recursively flattens nested columns in JSON data.
  json.infer-schema.flatten-nested-columns.enable: true
  # Uses the value of the databaseName field as the database name and the value of the tableName field as the table name.
  json.decode.parser-table-id.fields: databaseName,tableName
  # (Optional) Skips the first 100 parsing exceptions. If the number of exceptions exceeds 100, the job fails.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  # The metastore type. The value is fixed to rest.
  catalog.properties.metastore: rest
  # The token provider. The value is fixed to dlf.
  catalog.properties.token.provider: dlf
  # The URI to access the DLF Rest Catalog Server. The format is http://[region-id]-vpc.dlf.aliyuncs.com. For example, http://cn-hangzhou-vpc.dlf.aliyuncs.com.
  catalog.properties.uri: dlf_uri
  # The name of the DLF catalog.
  catalog.properties.warehouse: your_warehouse
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true

# Adds primary key information to the table. 
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
    
# Writes data from ods.inventory, ods.customer, and ods.user to the test_database.inventory, test_database.customer, and test_database.user tables, respectively.
route:
  - source-table: ods.inventory
    sink-table: test_database.inventory
  - source-table: ods.customer
    sink-table: test_database.customer
  - source-table: ods.user
    sink-table: test_database.user   
    
pipeline:
  # (Optional) Records dirty data that causes processing exceptions in the logs.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

For details on schema inference and evolution policies for JSON-formatted Kafka source tables, see Schema Parsing and Evolution Policies.

Use cases

The following sections describe job configurations for common use cases. For more detailed configurations, see Data Ingestion Kafka Connector.

Resolve field name conflicts

An Apache Kafka message consists of a Key and a Value. You can define the format for each part by setting key.format and value.format. The final schema is a combination of all fields from both parts.

If there are conflicting field names between the Key and Value parts, use key.fields-prefix and value.fields-prefix to add prefixes to the field names.

For example, if the Key part of a Kafka message contains id and name fields, and the Value part contains id and price fields, the following job configuration produces a schema with the fields key_id, key_name, val_id, and val_price.

source:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  key.format: json
  # Adds a prefix to field names in the Key.
  key.fields-prefix: key_
  # Adds a prefix to field names in the Value.
  value.fields-prefix: val_
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100

# Writes all data from the test_topic topic to the test_database.test_topic table.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Read metadata

Use the metadata.list configuration to read and pass additional Kafka message metadata. Metadata columns added in metadata.list can be used directly in the transform module. For a list of supported metadata, see Available metadata columns.

The following configuration adds the partition and offset metadata to the data. It then uses this metadata in the transform module to filter data where the partition is greater than 1 and the offset is greater than 100.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Adds metadata columns.
  metadata.list: partition,offset
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100
  
transform:
  - source-table: \.*.\.*
    filter: '`partition` > 1 and `offset` > 100'
    
# Writes all data from the test_topic topic to the test_database.test_topic table.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Handle dirty data

Log data may contain incorrectly formatted dirty data, which can cause the job to fail and restart repeatedly. Flink CDC data ingestion supports ignoring parsing errors and collecting the data that fails to parse. For more information, see Dirty Data Collection.

The following job writes dirty data that fails to parse to a log file. The job fails if more than 100 parsing errors occur.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100
  
# Writes all data from the test_topic topic to the test_database.test_topic table.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

If you do not need to collect dirty data and want to prevent job failures due to dirty data, you can use json.ignore-parse-errors, debezium-json.ignore-parse-errors, or canal-json.ignore-parse-errors to ignore parsing errors directly.

Specify a TableID

Debezium JSON and Canal JSON formats are fixed, with specific fields used to store the TableID. For generic JSON data, which has no fixed format, the topic name serves as the TableID by default. To use specific fields from the data as the TableID, configure json.decode.parser-table-id.fields. For example, given the JSON data {"col0":"a", "col1":"b", "col2":"c"}, different configurations generate the following TableIDs:

Configuration

TableID

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

The following job configuration concatenates the db and table columns from the data to create the TableID.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Uses the db and table fields from the data as the TableID.
  json.decode.parser-table-id.fields: db,table
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Infer data types

The data types for log data are inferred by parsing the data. For detailed information on type inference, see Schema Parsing and Evolution Policies. The following sections describe how to adjust configurations to control type parsing in common scenarios.

(General) Set all field types to string

If downstream processing and storage do not require specific field types, you can enable json.infer-schema.primitive-as-string, debezium-json.infer-schema.primitive-as-string, or canal-json.infer-schema.primitive-as-string. This skips type inference and sets all field types to String.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Sets the type for all fields in the JSON data to String.
  json.infer-schema.primitive-as-string: true
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100

# Writes all data from the test_topic topic to the test_database.test_topic table.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

(General) Specify an initial schema

In some scenarios, you may need to specify an initial schema, such as when writing Kafka data to a pre-existing downstream table. You can specify an initial schema by adding the scan.value.initial-schemas.ddls parameter.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Uses the db and table fields from the data as the TableID.
  json.decode.parser-table-id.fields: db,table
  # Sets the initial schema.
  scan.value.initial-schemas.ddls: |
    CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10));
    CREATE TABLE db1.t2 (id BIGINT);
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

The preceding configuration specifies an initial type of BIGINT for the id field and VARCHAR(10) for the name field in the db1.t1 table. It also specifies an initial type of BIGINT for the id field in the db1.t2 table.

(JSON) Set fixed field types

For JSON data, field types are inferred by parsing the JSON node types. Sometimes the inferred types may not meet your expectations. To resolve this, use the json.infer-schema.fixed-types configuration to specify types for certain fields.

The following job configuration specifies the id field type as BIGINT and the name field type as VARCHAR(10).

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Sets a fixed type for specific fields.
  json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
  # Required for versions 11.5 and earlier.
  scan.max.pre.fetch.records: 0
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100

# Writes all data from the test_topic topic to the test_database.test_topic table.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic
    
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

(Canal JSON) Specify the inference source

Canal JSON data contains more information than standard JSON data. If sqlType or mysqlType fields exist in the Canal JSON data, you can use this information to parse more precise data types.

  • Parsing schema from sqlType

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: canal-json
  # Infers schema from the sqlType field.
  canal-json.infer-schema.strategy: SQL_TYPE
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger
  • Parsing schema from mysqlType

    source:
      type: kafka
      name: Kafka source
      properties.bootstrap.servers: localhost:9092
      topic: test_topic
      properties.group.id: test_group
      scan.startup.mode: earliest-offset
      value.format: canal-json
      # Infers schema from the mysqlType field.
      canal-json.infer-schema.strategy: MYSQL_TYPE
      # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
      schema.inference.strategy: continuous
      # Ignores errors during data parsing.
      ingestion.ignore-errors: true
      # If 100 data parsing errors occur, the job fails.
      ingestion.error-tolerance.max-count: 100
      
    sink:
      type: paimon
      catalog.properties.metastore: rest
      catalog.properties.uri: dlf_uri
      catalog.properties.warehouse: your_warehouse
      catalog.properties.token.provider: dlf
      # (Optional) Enables deletion vectors to improve read performance.
      table.properties.deletion-vectors.enabled: true
      
    pipeline:
      # Enables dirty data collection. Dirty data is written to a log file.
      dirty-data.collector:
        name: Logger Dirty Data Collector
        type: logger

Use a static schema

If the schema of the data in your topic is fixed, set schema.inference.strategy to static. The data ingestion job then performs schema inference only once at startup and does not parse the schema of subsequent data.

source:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Sets the schema inference strategy to static. The schema is inferred only once when the job starts.
  schema.inference.strategy: static
  # Tries to consume 20 records from each partition to infer the schema.
  scan.max.pre.fetch.records: 20
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100
  
# Writes all data from the test_topic topic to the test_database.test_topic table.
route:
  - source-table: test_topic
    sink-table: test_database.test_topic
    
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

You can also add the scan.value.initial-schemas.ddls parameter to specify an initial schema and skip schema inference for certain tables. The following example specifies initial schemas for the db1.t1 and db1.t2 tables.

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Uses the db and table fields from the data as the TableID.
  json.decode.parser-table-id.fields: db,table
  # Sets the initial schema.
  scan.value.initial-schemas.ddls: | 
    CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10));
    CREATE TABLE db1.t2 (id BIGINT);
  # Sets the schema inference strategy to static. The schema is inferred only once when the job starts.
  schema.inference.strategy: static
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Accelerate Kafka log data parsing

In addition to general Kafka connector acceleration configurations, the Flink CDC data ingestion job offers specific settings to speed up parsing based on your use case.

  1. Schema parsing can be time-consuming. If your downstream application only requires the String type, you can enable json.infer-schema.primitive-as-string, debezium-json.infer-schema.primitive-as-string, or canal-json.infer-schema.primitive-as-string. This skips type inference and sets all field types to String, which accelerates parsing.

  2. For Canal JSON data, you can use canal-json.database.include and canal-json.table.include to filter out data from unnecessary tables.

  3. If the data schema will not change or if you do not need schema evolution, you can change schema.inference.strategy to static. This performs schema inference only once when the job starts.

SLS log ingestion into a data lake

Log Service (SLS) is a one-stop service for log data. Flink CDC data ingestion lets you quickly ingest log data from SLS into a data lake in real time. It automatically performs schema inference and supports schema evolution.

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger
  • TableID: By default, the project and logstore are concatenated to form the TableID. In the job above, the TableID would be test_pj.test_log.

  • Data type: The SLS connector treats all fields in each log entry as the String type by default.

  • Schema evolution: Schema evolution is currently limited to adding new columns. New columns are appended to the end of the schema by default.

Use cases

The following sections describe job configurations for common use cases. For more detailed configurations, see Data Ingestion SLS Connector.

Read metadata

You can use the metadata.list configuration to read and pass additional SLS metadata. Metadata columns added in metadata.list can be used directly in the transform module. For a list of supported metadata, see metadata.list.

Note that columns added with metadata.list are not automatically included in the output data. To write these metadata columns to the sink, you must declare them in the projection section of the transform module.

The following configuration adds the __timestamp__ and __tag__ metadata to the data and writes them to the sink. It also uses these metadata fields in the transform module to filter data where the timestamp is greater than 1772181154 and the tag is "test".

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Adds metadata columns.
  metadata.list: __timestamp__,__tag__
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100
  
transform:
  - source-table: \.*.\.*
    projection: \*, __timestamp__ as timestamp_col, __tag__ as tag_col
    filter: '`__timestamp__` > 1772181154 and `__tag__` = "test"'

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Handle dirty data

Log data may contain incorrectly formatted dirty data, which can cause the job to fail and restart repeatedly. Flink CDC data ingestion supports ignoring parsing errors and collecting the data that caused the errors. For more information, see Dirty Data Collection.

The following job writes dirty data that fails to parse to a log file. The job fails if more than 100 parsing errors occur.

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Specify a TableID

To use specific fields from the data as the TableID, configure decode.table-id.fields. For example, given the log data {"col0":"a", "col1":"b", "col2":"c"}, different configurations generate the following TableIDs:

Configuration

TableID

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

The following job configuration concatenates the db and table columns from the data to create the TableID.

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Uses the db and table fields from the data as the TableID.
  decode.table-id.fields: db,table
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Specify field types

The data ingestion SLS connector treats all fields in the data as the String type by default. If you need to specify the type for certain fields, use the fixed-types configuration option.

The following job specifies the id field type as BIGINT and the name field type as VARCHAR(10).

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Specifies the id field type as BIGINT and the name field type as VARCHAR(10).
  fixed-types: id BIGINT, name VARCHAR(10)
  # Sets the schema inference strategy to continuous. This strategy detects the schema of each message and synchronizes schema changes.
  schema.inference.strategy: continuous
  # Ignores errors during data parsing.
  ingestion.ignore-errors: true
  # If 100 data parsing errors occur, the job fails.
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (Optional) Enables deletion vectors to improve read performance.
  table.properties.deletion-vectors.enabled: true
  
pipeline:
  # Enables dirty data collection. Dirty data is written to a log file.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger