Data Ingestion
On this page
Engine support: VERA 4.6 (Flink 1.20).
What's New in VERA 4.6
Data ingestion has exited public beta and is now generally available.
Merging columns into a lake
Merge multiple fields from an upstream JSON source, with different names or casing, into a single target column. This supports regular-expression matching, case normalisation, and custom mappings.
Append-only partitioned writes
The Paimon sink can now write to partitioned tables that have no primary key, for append-only use cases. You no longer need to include the partition key in the primary key.
Clearing PrimaryKey or PartitionKey
As of VERA 4.6, you can clear the primary key inferred from the upstream and convert the table to one without a primary key, by passing an empty value:
1transform:
2 - source-table: db.tbl
3 primary-keys: '' # Clears the PRIMARY KEY setting for the db.tbl tableTable-name routing with regular expressions
Define more complex table-name routing logic using regular expressions.
Variant full-link support
CDC YAML supports accessing VARIANT-type fields, converting them, and writing them to a Paimon sink.
Kafka source enhancements
The Kafka source can split a single message by fields into multiple records, written to different target tables (field routing), and supports custom partitioners.
Sink enhancements
Fluss sink: supports writing database CDC, with native schema evolution.
Paimon sink: the commit node's concurrency can now be configured independently.
Iceberg sink: supports built-in catalog references and automatic retrieval of connection information (for example URL and account credentials), so configuration can be reused.
New built-in functions
MD5 is now available as a hash function in data ingestion transforms.
parse_json converts a JSON string into a VARIANT value.
Why YAML for CDC?
- Simplified job management – Declarative, human‑readable configuration.
- Reusability & consistency – Template values, reuse across envs.
- CI/CD‑friendly – Store in Git, review, promote, rollback.
- Environment separation – Swap credentials/topics/URIs per env.
- Faster onboarding – No Flink internals required.
- Tooling compatibility – Validate/lint/test YAML.
- Separation of concerns – Data flow vs. runtime/platform config.
Quick start (UI)
- Go to Data Ingestion → New Draft and select Blank Draft.
- Name the draft and pick your Engine Version (match to what the job was tested with).
- Paste a YAML CDC config in the Preview panel and click OK.
You can also create drafts from the SQL Editor or import files from your repo.
YAML schema overview
At minimum, a CDC YAML config contains source and sink sections. A job can include one source and one sink per file; compose multiple files for multiple flows.
1source:
2 type: <connector> # e.g., mysql, postgres, oracle, sqlserver, kafka
3 name: <human name>
4 hostname: <host or service>
5 port: <int>
6 username: ${secret_values.mysqlusername}
7 password: ${secret_values.mysqlpassword}
8 database: <db-name> # optional depending on connector
9 tables: <regex or list> # e.g., "mysql\.\.*" or [ db.schema.table1, db.schema.table2 ]
10 server-id: <range> # connector-specific; example for MySQL
11 snapshot.mode: initial # connector-specific snapshot policy
12
13sink:
14 type: <connector> # e.g., mysql, postgres, kafka, iceberg, hudi, jdbc
15 name: <human name>
16 hostname: <host-or-broker>
17 port: <int>
18 username: <user>
19 password: <pass>
20 database: <db>
21 table: <table or pattern>
22 upsert-key: <col or list> # for upsert sinksSecrets & variables.
Use ${…} expressions (e.g., ${secret_values.mysqlpassword}) to reference values injected at deploy time. Treat credentials as secrets. Do not hardcode.
Minimal example – MySQL → MySQL
1source:
2 type: mysql
3 name: Database A to Data warehouse
4 hostname: mysql-src
5 port: 3306
6 username: ${secret_values.mysqlusername}
7 password: ${secret_values.mysqlpassword}
8 tables: mysql\.\.*
9
10sink:
11 type: mysql
12 name: Database B to Data warehouse
13 hostname: mysql-dst
14 port: 3306
15 username: root
16 password: passCommon fields and patterns
Table selection
- Single table:
tables: mydb.public.users - Multiple tables:
1tables:
2 - mydb.public.users
3 - mydb.public.orders- Regex pattern:
tables: mydb\.public\..*(Escape dots in YAML strings.)
Primary/Upsert keys
For sinks that support upserts, specify a key:
1sink:
2 type: mysql
3 table: dw.users
4 upsert-key: idParallelism & checkpoints (runtime)
Runtime parameters are set on the deployment and can be overridden per job if supported:
1runtime:
2 parallelism: 4
3 checkpoint-interval: 60s
4 restart-strategy: fixed-delayError handling
1on-error:
2 drop: false # default; fail the job on deserialization errors
3 dead-letter: # optional DLQ
4 type: kafka
5 topic: cdc-dlqExact runtime/error keys vary by connector; prefer the connector’s reference.