Sink Connectors

Sink Connectors: Loading BookNest Events into a Database

A sink connector is a managed consumer group: tasks take partitions, write each batch, and Connect commits offsets after the writes succeed. The JDBC sink needs a schema per record to build INSERT and CREATE TABLE; the shop's events are plain JSON, so the converter gets one from the config.

listings/order-event.schema.json: the order event in Connect's schema notationJSON
{"type": "struct", "name": "OrderEvent", "fields": [
  {"field": "event_id", "type": "int64"}, {"field": "ts", "type": "string"},
  {"field": "type", "type": "string"}, {"field": "order_id", "type": "int64"},
  {"field": "customer_id", "type": "int64", "optional": true},
  {"field": "total", "type": "double", "optional": true}]}
listings/orders-sink.json: upsert every order event into PostgreSQLJSON
{
  "connector.class": "io.aiven.connect.jdbc.JdbcSinkConnector",
  "topics": "booknest.order-events", "tasks.max": "3",
  "connection.url": "jdbc:postgresql://l3-pg:5432/booknest",
  "connection.user": "booknest", "connection.password": "booknest",
  "table.name.format": "order_events", "auto.create": "true",
  "insert.mode": "upsert", "pk.mode": "record_value", "pk.fields": "event_id",
  "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  "value.converter.schemas.enable": "true",
  "transforms": "ts,rename",
  "transforms.ts.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value",
  "transforms.ts.field": "ts", "transforms.ts.target.type": "Timestamp",
  "transforms.ts.format": "yyyy-MM-dd'T'HH:mm:ssX",
  "transforms.rename.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
  "transforms.rename.renames": "type:event_type"
}
listings/l0685_sink.sh: add the schema, start the sink, and time the loadShell
# Add the event schema to the sink's config, start the sink, and time the load into PostgreSQL
C=localhost:33083/connectors/booknest-orders-sink
jq --arg s "$(jq -c . listings/order-event.schema.json)" \
  '. + {"value.converter.schema.content": $s}' listings/orders-sink.json > sink.json
start=$(date +%s)
curl -s -X PUT -H 'Content-Type: application/json' -d @sink.json $C/config >/dev/null
q() { docker exec l3-pg psql -U booknest -Atc "$1" 2>/dev/null; }
until [ "$(q 'SELECT count(*) FROM order_events')" = 390737 ]; do sleep 1; done
echo "390,737 rows in $(( $(date +%s) - start )) s"
q "SELECT string_agg(column_name || ' ' || udt_name, ', ' ORDER BY ordinal_position)
   FROM information_schema.columns WHERE table_name = 'order_events'"
Output
390,737 rows in 18 s
event_id int8, ts timestamp, event_type text, order_id int8, customer_id int8, total float8

Three tasks, one per partition, loaded the events in under half a minute on this shared four-CPU machine. Two transforms (Single Message Transforms) made ts a timestamp and renamed type. Upsert on event_id makes replays after a crash overwrite rows instead of duplicating them. Analytics-scale landing belongs in Iceberg 129 , through the sink connector of Streaming into the Lakehouse.