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: loggerAssume 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: loggerFor 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: loggerRead 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: loggerHandle 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: loggerIf 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: loggerInfer 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: loggerThe 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: loggerParsing schema from
mysqlTypesource: 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: loggerYou 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: loggerAccelerate 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.
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, orcanal-json.infer-schema.primitive-as-string. This skips type inference and sets all field types to String, which accelerates parsing.For Canal JSON data, you can use
canal-json.database.includeandcanal-json.table.includeto filter out data from unnecessary tables.If the data schema will not change or if you do not need schema evolution, you can change
schema.inference.strategytostatic. 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: loggerTableID: By default, the
projectandlogstoreare concatenated to form the TableID. In the job above, the TableID would betest_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: loggerHandle 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: loggerSpecify 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: loggerSpecify 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