ClickZetta Spark & Flink Connector
Read references/spark.md for the Spark Connector and references/flink.md for the Flink Write Connector.
Key Constraints (Required Reading)
| Constraint |
Spark Connector |
Flink Connector |
| Primary-key table writes |
Not supported |
Supported in igs-dynamic-table mode |
| Partial-column writes |
Not supported; all columns must be written |
Supported |
| CDC (UPDATE/DELETE) |
Not supported; append only |
Supported in igs-dynamic-table mode |
| Spark version |
3.4.0+ |
N/A |
| Flink version |
N/A |
1.14, 1.15, 1.17, 1.18 |
Spark Connector Quick Example
// Write
df.write.format("clickzetta")
.option("endpoint", "your_instance.cn-shanghai-alicloud.api.clickzetta.com")
.option("username", sys.env("CZ_USERNAME"))
.option("password", sys.env("CZ_PASSWORD"))
.option("workspace", "your_workspace")
.option("virtualCluster", "default")
.option("schema", "public")
.option("table", "orders")
.mode("append")
.save()
// Read
val df = spark.read.format("clickzetta")
.option("endpoint", "your_instance.cn-shanghai-alicloud.api.clickzetta.com")
.option("username", sys.env("CZ_USERNAME"))
.option("password", sys.env("CZ_PASSWORD"))
.option("workspace", "your_workspace")
.option("virtualCluster", "default")
.option("schema", "public")
.option("table", "orders")
.load()
Flink Connector Quick Example
-- CDC mode (supports INSERT/UPDATE/DELETE; the target table must have a primary key)
CREATE TABLE lakehouse_sink (
order_id INT,
status STRING,
amount DOUBLE,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'igs-dynamic-table',
'curl' = 'jdbc:clickzetta://your_instance.cn-shanghai-alicloud.api.clickzetta.com/default?username=user&password=***&schema=public',
'schema-name' = 'public',
'table-name' = 'orders',
'sink.parallelism' = '1',
'properties' = 'authentication:true'
);
INSERT INTO lakehouse_sink SELECT order_id, status, amount FROM source_table;
Selection Guide
| Scenario |
Recommended option |
| Spark ETL batch writes to a non-primary-key table |
Spark Connector |
| Flink real-time stream writes to a non-primary-key table |
Flink igs-dynamic-table-append-only |
| Flink CDC synchronization to a primary-key table, including UPDATE/DELETE |
Flink igs-dynamic-table |
| High-frequency real-time writes from a Java application |
Java SDK RealtimeStream |
1---2name: clickzetta-spark-flink-connector3description: Write data to ClickZetta Lakehouse using the Spark Connector or Flink Write Connector. Covers Spark DataFrame read/write configuration (Maven dependencies, connection parameters, read/write code), Flink Table API writes (CDC mode igs-dynamic-table, append-only mode igs-dynamic-table-append-only), checkpoint configuration, buffer/flush tuning, and key constraints such as primary-key table limitations. Trigger when the user says "Spark Connector", "Flink Connector", "Spark writes to Lakehouse", "Flink writes to Lakehouse", "spark-clickzetta", "igs-flink-connector", "Spark DataFrame write", "Flink CDC write", "Flink sink", or "spark.read.format clickzetta". Keywords: Spark, Flink, DataFrame, connector, read, write, CDC, igs-dynamic-table4---56# ClickZetta Spark & Flink Connector78Read [references/spark.md](references/spark.md) for the Spark Connector and [references/flink.md](references/flink.md) for the Flink Write Connector.910---1112## Key Constraints (Required Reading)1314| Constraint | Spark Connector | Flink Connector |15|---|---|---|16| Primary-key table writes | Not supported | Supported in `igs-dynamic-table` mode |17| Partial-column writes | Not supported; all columns must be written | Supported |18| CDC (UPDATE/DELETE) | Not supported; append only | Supported in `igs-dynamic-table` mode |19| Spark version | 3.4.0+ | N/A |20| Flink version | N/A | 1.14, 1.15, 1.17, 1.18 |2122---2324## Spark Connector Quick Example2526```scala27// Write28df.write.format("clickzetta")29 .option("endpoint", "your_instance.cn-shanghai-alicloud.api.clickzetta.com")30 .option("username", sys.env("CZ_USERNAME"))31 .option("password", sys.env("CZ_PASSWORD"))32 .option("workspace", "your_workspace")33 .option("virtualCluster", "default")34 .option("schema", "public")35 .option("table", "orders")36 .mode("append")37 .save()3839// Read40val df = spark.read.format("clickzetta")41 .option("endpoint", "your_instance.cn-shanghai-alicloud.api.clickzetta.com")42 .option("username", sys.env("CZ_USERNAME"))43 .option("password", sys.env("CZ_PASSWORD"))44 .option("workspace", "your_workspace")45 .option("virtualCluster", "default")46 .option("schema", "public")47 .option("table", "orders")48 .load()49```5051---5253## Flink Connector Quick Example5455```sql56-- CDC mode (supports INSERT/UPDATE/DELETE; the target table must have a primary key)57CREATE TABLE lakehouse_sink (58 order_id INT,59 status STRING,60 amount DOUBLE,61 PRIMARY KEY (order_id) NOT ENFORCED62) WITH (63 'connector' = 'igs-dynamic-table',64 'curl' = 'jdbc:clickzetta://your_instance.cn-shanghai-alicloud.api.clickzetta.com/default?username=user&password=***&schema=public',65 'schema-name' = 'public',66 'table-name' = 'orders',67 'sink.parallelism' = '1',68 'properties' = 'authentication:true'69);7071INSERT INTO lakehouse_sink SELECT order_id, status, amount FROM source_table;72```7374---7576## Selection Guide7778| Scenario | Recommended option |79|---|---|80| Spark ETL batch writes to a non-primary-key table | Spark Connector |81| Flink real-time stream writes to a non-primary-key table | Flink `igs-dynamic-table-append-only` |82| Flink CDC synchronization to a primary-key table, including UPDATE/DELETE | Flink `igs-dynamic-table` |83| High-frequency real-time writes from a Java application | Java SDK RealtimeStream |