บทที่ 7: Ingestion
2 min readBatch Ingestion Patterns
| รูปแบบ | คำอธิบาย |
|---|---|
| Full Snapshot | ดึงสถานะทั้งหมดทุกครั้ง — ง่ายแต่ใช้ bandwidth เยอะ |
| Incremental | ดึงเฉพาะการเปลี่ยนแปลง — ลด network/storage แต่ซับซ้อนกว่า |
| ETL (Extract → Transform → Load) | Transform ก่อนโหลด — ระบบเก่า, on-premises (destination มีพลังประมวลผลจำกัด) |
| ELT (Extract → Load → Transform) | โหลดก่อนแล้วค่อย transform — cloud data warehouse (ใช้พลังของ warehouse ช่วย transform) |
ตัด batch ได้ 2 วิธี: Time-interval (ตามช่วงเวลา เช่น รันทุกคืน — พบบ่อยใน ETL แบบดั้งเดิม) หรือ Size-based (ตามขนาดข้อมูลที่สะสม เช่น ก่อนเขียนลง object storage เป็นก้อนๆ)
Figure 7-4. Spectrum ของความถี่ ingestion ตั้งแต่ batch ช้าไปจนถึง real-time
Streaming Ingestion Patterns
- Real-time ที่แท้จริงไม่มี — ทุกระบบมี latency
- Near real-time: process ทีละ events หรือ micro-batches
| แพลตฟอร์ม | จุดเด่น |
|---|---|
| Kafka | Persistent log, replay, schema registry |
| Kinesis | Shard-based scaling, max 365 days retention |
| Pub/Sub | Push + pull subscriptions, managed |
Figure 7-6. ตัวอย่าง asynchronous event pipeline บน AWS — Kinesis (buffer) → Beam (enrich) → Kinesis → Firehose → S3
Key Challenges
- Schema Evolution: ใช้ schema registry + dead-letter queue
- Late-Arriving Data: ต้องตั้ง cutoff time
- Duplicates & Exactly-Once: streaming platforms มัก deliver at-least-once
- Throughput: ใช้ buffering และ managed services
- Bursty data: ต้องการ buffering เสมอ
- Time to Live (TTL): เวลาสูงสุดที่ event รอ ack ก่อนถูกทิ้ง — สั้นไปข้อความหายก่อนประมวลผล, ยาวไปข้อความค้างเยอะ (Kafka เก็บได้ไม่จำกัดถ้ามีพื้นที่พอ, Kinesis สูงสุด 365 วัน, Pub/Sub สูงสุด 7 วัน)
Figure 7-12. Event ที่ ingest สำเร็จ ("good") ส่งต่อไป consumer ส่วนที่ล้มเหลว ("bad") ถูกเก็บใน dead-letter queue
Tools & Technologies
| ประเภท | ตัวอย่าง |
|---|---|
| Direct DB Connection | JDBC/ODBC — ดึงข้อมูลผ่าน network connection ตรงจาก database |
| Managed Connectors | Fivetran, Airbyte, Matillion |
| Streaming | Kafka, Kinesis, Pulsar, Pub/Sub |
| CDC | Debezium, AWS DMS, Striim |
| File Transfer | SFTP, Snowball, Object storage |
| Webhook | "Reverse API" — ผู้ให้ข้อมูลเป็นฝ่าย call API ของเรา แทนที่เราจะ pull เอง |
Figure 7-13. Ingestion process อ่านจาก source database ด้วย JDBC แล้วเขียนเป็น object ลง object storage
Figure 7-15. สถาปัตยกรรม webhook ingestion ทั่วไปบน cloud services