PrepZone Logo
PrepZone

Kafka Connect: Sources, Sinks, and SMTs

Move data in and out of Kafka without custom consumers — JDBC source, S3 sink, and single message transforms.

Why this matters

  • VaultCommerce uses JDBC source for legacy MySQL reports and S3 sink for compliance archives.
  • Distributed Connect workers share connector tasks across the cluster for fault tolerance.
  • SMTs (Single Message Transforms) filter and reshape records in the pipeline.
  • Connect handles offset management for sources — different from consumer groups.
Postgres / S3Source
Kafka topics
Elasticsearch / S3Sink
Source connectors ingest external data. Sink connectors export to warehouses, indexes, or object storage.

Source and sink connectors

Source: JDBC polls updated_at column → Kafka topic. Sink: S3 connector batches by size/time → Parquet files. VaultCommerce runs 3 Connect workers with group.id=vaultcommerce-connect for HA.

Key points

  • Connector — plugin configuration (Debezium, JDBC, S3, Elasticsearch)
  • Task — parallel unit of work; partition of source table or topic
  • Worker — JVM running connector tasks in distributed mode
  • SMT — in-flight record transform without custom code
  • Converter — JSON, Avro, or ByteArray between Connect and Kafka format

SMTs and error handling

ExtractField pulls nested JSON; Filter drops test tenants. errors.tolerance=all with errors.deadletterqueue.topic.name routes poison rows. VaultCommerce never uses tolerance=all on financial connectors.

VaultCommerce rollout checklist

Before promoting changes that touch the VaultCommerce order and payment event backbone, run the staging KRaft cluster (Kafka 3.7+, Schema Registry 7.x) through a 10k events/min soak test. Compare producer request latency p99 and consumer lag per group against the pre-deploy baseline. Connect for standardized system integration. Document the change in the internal topic registry, attach Grafana screenshots to the change ticket, and keep an engineer on lag dashboards for 30 minutes after production rollout — roll back the service release before altering broker-level settings if lag or under-replicated partitions spike.

Java
{
  "name": "vaultcommerce-s3-archive-sink",
  "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "topics": "vaultcommerce.orders.placed.v1",
    "s3.bucket.name": "vaultcommerce-kafka-archive",
    "s3.region": "us-east-1",
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
    "path.format": "'year'=YYYY/'month'=MM/'day'=dd",
    "flush.size": "10000",
    "tasks.max": "4"
  }
}

Quick recall

Everything you need if you only revisit this box.

  1. Connect for standardized system integration.
  2. Distributed workers for HA and task parallelism.
  3. SMTs for light transforms; not a replacement for stream processing.

Test yourself

Answer these before moving on — recall is what makes it stick.