JDBC Source Connector

Source Connectors and the JDBC Source Example

A JDBC source polls tables with SQL. Its mode picks how to find changes: bulk (the whole table), incrementing (a growing ID), timestamp (a last-modified column) or timestamp+incrementing, which loses no row that shares a timestamp. The last (updated_at, id) seen is the source offset.

listings/catalog-source.json: poll BookNest's catalog table for new and changed booksJSON
{
  "connector.class": "io.aiven.connect.jdbc.JdbcSourceConnector",
  "connection.url": "jdbc:postgresql://l3-pg:5432/booknest",
  "connection.user": "booknest", "connection.password": "booknest",
  "table.whitelist": "catalog", "topic.prefix": "booknest.pg.",
  "mode": "timestamp+incrementing",
  "timestamp.column.name": "updated_at", "incrementing.column.name": "id",
  "numeric.mapping": "best_fit", "poll.interval.ms": "2000", "tasks.max": "1"
}
listings/l0684_source.sh: start Connect, create the source, change a price, read the topicShell
# Start Connect, create (or update) the source with PUT, change a price, read the topic
docker compose -f compose/connect.yaml up -d --wait 2>/dev/null
C=localhost:33083/connectors/booknest-catalog-source
curl -s -X PUT -H 'Content-Type: application/json' -d @listings/catalog-source.json \
  $C/config >/dev/null && sleep 10
docker exec l3-pg psql -U booknest -qc \
  "UPDATE catalog SET price = 34.50, updated_at = now() WHERE id = 2" && sleep 5
docker exec l3-kafka /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server l3-kafka:9092 \
  --topic booknest.pg.catalog --from-beginning --timeout-ms 5000 2>/dev/null |
  jq -c '{id, title, price, updated_at}'
Output
{"id":1,"title":"The Quiet Harbor","price":14.99,"updated_at":1790917777288}
{"id":2,"title":"Patterns of the Deep Web","price":39.5,"updated_at":1790917777288}
...
{"id":2,"title":"Patterns of the Deep Web","price":34.5,"updated_at":1790917835423}

The first poll copied the six books into booknest.pg.catalog (prefix plus table); the new price followed two seconds later. numeric.mapping=best_fit sends numeric(6,2) as a double, not a Base64 Decimal. Polling has limits: every change must bump updated_at, updates between polls merge, each poll queries the database, and deletes are invisible.