ข้ามไปยังเนื้อหา

Kafka Connect

Kafka Connect คือ framework สำเร็จรูปสำหรับย้ายข้อมูล เข้าและออก จาก Kafka ด้วยการตั้งค่าแทนการเขียนโค้ด — ชี้ source connector ไปที่ database หรือไฟล์เพื่อดึงข้อมูลเข้ามา หรือชี้ sink connector ไปที่ search index หรือ warehouse เพื่อส่งข้อมูลออกไป แล้ว Connect จะจัดการเรื่อง scaling, offset และการ restart ให้เอง

งาน “เกี่ยวกับ Kafka” จำนวนมากไม่ใช่ business logic เลย — แต่เป็นงานท่อ: อ่าน row จาก Postgres เข้า topic หรือเขียน topic ลง Elasticsearch การเขียนท่อพวกนี้เองหมายถึงต้องแก้ปัญหา offset tracking, retry และ parallelism ซ้ำทุกครั้ง Connect แพ็กงานพวกนี้เป็น connector ที่ใช้ซ้ำได้:

  • source connector — ระบบภายนอก → Kafka (database, ไฟล์, stream แบบ change-data-capture)
  • sink connector — Kafka → ระบบภายนอก (database, search index, data warehouse, object storage)

เพราะ connector เป็นแค่ config จึงมีให้ใช้อยู่แล้วเป็นร้อย ๆ ตัว คุณแทบไม่ต้องเขียนเอง แค่ตั้งค่าเท่านั้น

// A JDBC source connector: poll a Postgres table into a topic
{
"name": "orders-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"connection.url": "jdbc:postgresql://db:5432/shop",
"table.whitelist": "orders", // which table to pull
"mode": "incrementing", // track new rows by a column
"incrementing.column.name": "id",
"topic.prefix": "pg-", // -> topic "pg-orders"
"tasks.max": "3" // up to 3 parallel tasks
}
}

มีสามชิ้นส่วนที่ทำงานร่วมกัน:

  • worker — process JVM ที่รัน connector หนึ่ง worker หรือมากกว่ารวมกันเป็น Connect cluster
  • task — หน่วยของ parallelism connector แบ่งงานเป็น task ได้ถึง tasks.max ตัว แล้ว worker กระจาย task ไปทั่ว cluster (เช่น หนึ่ง task ต่อหนึ่ง table หรือต่อช่วง partition)
  • converter — แปลงระหว่างรูปแบบ record ภายในของ Connect กับ bytes บน topic ตรงนี้เองที่ Schema Registry เสียบเข้ามา: AvroConverter serialize พร้อม schema id ส่วน JsonConverter เขียนเป็น JSON ธรรมดา
flowchart LR
  db["Postgres"] -->|"source connector"| worker["Connect worker + tasks"]
  worker -->|"converter"| topic["Topic: pg-orders"]
  topic -->|"converter"| worker2["Connect worker + tasks"]
  worker2 -->|"sink connector"| es["Elasticsearch"]
Connect moves data between external systems and Kafka

Connect รันได้สองโหมด:

  • standalone — worker เดียว config อยู่ในไฟล์ local ง่าย ไม่มี fault tolerance เหมาะกับ demo บนเครื่องหรืองาน ingest ไฟล์ครั้งเดียวจบ
  • distributed — หลาย worker ทำงานประสานกันเหมือน consumer group โดยเก็บ config ของ connector และ offset ไว้ ใน topic ของ Kafka ถ้า worker ตัวหนึ่งตาย task ของตัวเองจะ rebalance ไปยังตัวที่รอด นี่คือสิ่งที่คุณรันบน production และจัดการผ่าน REST API:
Terminal window
# Deploy the connector above to a distributed Connect cluster
curl -X POST http://localhost:8083/connectors \
-H "Content-Type: application/json" \
-d @orders-source.json

หยิบ Connect มาใช้เมื่องานคือ ย้ายข้อมูลโดยไม่มี logic พิเศษ — ingest table, mirror topic ไป S3, ป้อน search index คุณได้ offset management, การ restart และ scaling มาฟรี ให้เขียน producer หรือ consumer แทนเมื่อคุณต้องการ business logic จริง ๆ: การ enrich กับอีก service, การ branch ที่ซับซ้อน หรือ transformation ที่โตเกิน converter ธรรมดา ถ้าคุณเริ่มอยาก transform และ join stream นั่นคืองานของบทถัดไป (Kafka Streams) ไม่ใช่งานของ Connect

source connector กับ sink connector ต่างกันอย่างไร
ใน Kafka Connect task คืออะไร
ทำไมต้องรัน Connect แบบ distributed บน production
เมื่อไรควรเขียน producer หรือ consumer แทนการใช้ Connect