streaming/load_events.sh recreates Apache Kafka and Managed Cloud Kafka's topic: three partitions, key order_id, and the 390,737 JSON lines of order_events.jsonl as values. The landing table is created first in Spark 129 (order_events.sql: event_id BIGINT, ts TIMESTAMP, type STRING, order_id BIGINT, customer_id BIGINT, total DECIMAL(10,2), partitioned by months(ts)), so you choose the types, not the first JSON record. The configuration:
{
"name": "booknest-events-sink",
"config": {
"connector.class": "org.apache.iceberg.connect.IcebergSinkConnector",
"tasks.max": "3",
"topics": "booknest.order-events",
"iceberg.tables": "booknest.order_events",
"iceberg.control.commit.interval-ms": "15000",
"iceberg.catalog.type": "rest",
"iceberg.catalog.uri": "http://l1-iceberg-rest:8181",
"iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",
"iceberg.catalog.s3.endpoint": "http://l1-minio:9000",
"iceberg.catalog.s3.path-style-access": "true",
"iceberg.catalog.s3.access-key-id": "booknest-admin",
"iceberg.catalog.s3.secret-access-key": "booknest-secret-2026",
"iceberg.catalog.client.region": "us-east-1"
}
}The worker's JsonConverter (schemas.enable=false) turns each value into a map. start_sink.sh posts the file to the REST API, and commits.sh lists the table's snapshots from the catalog:
{"name":"booknest-events-sink","tasks_max":"3"}
connector RUNNING
task 0 RUNNING
...
time op rows files total_rows total_files
02:10:01 append 390737 54 390737 54The backlog arrived as one snapshot of 54 files: three tasks times 18 monthly partitions. The interval defaults to five minutes. iceberg.tables.dynamic-enabled with route-field fans records out to tables named in a field, and auto-create-enabled and evolve-schema-enabled let the sink create tables and add columns, convenient for raw landing zones and risky for curated ones.