การซิงก์ข้อมูลจาก BigQuery เข้าระบบ
ภาพรวม
ลูกค้าองค์กรจำนวนมากเก็บข้อมูลลูกค้าไว้ใน Google BigQuery เช่น ยอดซื้อสะสม กลุ่มลูกค้า หรือคะแนน RFM
ระบบจึงเปิดให้เชื่อมต่อ BigQuery ขององค์กรเข้ามา แล้วซิงก์ค่าเหล่านั้นลง line_user.custom_attribute
เป็นประจำ ทำให้สามารถแบ่งกลุ่มผู้ใช้และตั้ง trigger จากข้อมูลฝั่ง data warehouse ได้
จุดเด่นของงานนี้คือ ตารางเวลาการทำงานถูกกำหนดโดยข้อมูลในฐานข้อมูล ไม่ใช่โดยโค้ด
แต่ละแถวใน bigquery_sync_config มี cron expression ของตัวเอง และใช้ credential ของ tenant ตัวเอง
Business Flow
- Profile
bigquery-syncเริ่มทำงาน โดยใช้เพียง PostgreSQL อย่างเดียว ไม่ต้องพึ่ง RabbitMQ หรือ Redis - โหลดทุกแถวใน
bigquery_sync_configที่เปิดใช้งานอยู่ แล้วลงทะเบียน cron job หนึ่งตัวต่อหนึ่งแถว ในชื่อรูปแบบbigquery-sync-<id>โดยใช้ค่าsync_cronของแถวนั้น หรือใช้ค่าเริ่มต้น0 6 * * *ซึ่งหมายถึงตี 6 ของทุกวัน - เมื่อถึงเวลาทำงาน ระบบจะ
- อ่าน config ใหม่จากฐานข้อมูลอีกครั้ง เผื่อมีการแก้ไขหลังลงทะเบียน cron ไปแล้ว
- สร้าง BigQuery client เฉพาะ tenant จาก
project_idและcredentials_json - Query ตารางสรุปตามที่ config ระบุไว้
- จับคู่และเขียนข้อมูล
- ค้นหา
line_userจากค่าในคอลัมน์ที่กำหนดในuser_id_column - แปลงชื่อคอลัมน์ของ BigQuery ไปเป็นคีย์ของ
custom_attributeตามที่ระบุในfield_mappings - เขียนข้อมูลแบบ jsonb merge ทีละชุด ชุดละ 100 แถว
- แถวที่เกิดข้อผิดพลาดจะถูกข้ามและนับเก็บไว้ โดยไม่ทำให้ทั้งงานล้มเหลว
- ค้นหา
- เขียนค่า
last_sync_atและlast_sync_resultซึ่งบรรจุจำนวนแถวที่อัปเดตสำเร็จและจำนวนข้อผิดพลาด กลับลงแถว config เพื่อให้ CMS แสดงสถานะการซิงก์ล่าสุดได้ - ตอนปิดระบบ ตัว scheduler จะรอให้ job ที่กำลังทำงานอยู่เสร็จก่อนจึงหยุด
ไฟล์และฟังก์ชันหลัก
| ไฟล์ | หน้าที่ |
|---|---|
internal/bigquerysync/service.go | Service.Start และ Service.Stop ควบคุมวงจรชีวิตของ cron scheduler ส่วนการทำงานจริงอยู่ที่ loadAndRegisterConfigs, registerCronForConfig, syncForConfig, createBigQueryClient, querySummaryTable, batchUpdateUsers, updateUserAttribute, updateSyncResult พร้อมค่าคงที่ขนาด batch 100 แถว และ cron เริ่มต้น 0 6 * * * |
internal/bigqueryx/bigqueryx.go | Wrapper ของ Google BigQuery client |
cmd/worker/main.go | runBigquerySync — จุดเริ่มของ profile bigquery-sync ซึ่งเปิด admin port ที่ :9107 |
งานนี้ไม่มีคิวที่เกี่ยวข้องเลย เป็น cron ล้วน ไม่ consume และไม่ publish ข้อความใดใน RabbitMQ
จุดเชื่อมต่อกับ Service อื่น
- ตาราง
bigquery_sync_config— เก็บline_oa_id,organization_id,project_id,credentials_json,summary_table,user_id_column,field_mappings,sync_cron, สถานะเปิดใช้งาน รวมถึงlast_sync_atและlast_sync_result - ตาราง
line_user— ปลายทางของการเขียนcustom_attributeแบบ jsonb merge - Google BigQuery — สร้าง client แยกตาม tenant จาก service account credential ที่เก็บในฐานข้อมูล
- cms-api-go — เจ้าของหน้าตั้งค่าการเชื่อมต่อและหน้าดูผลการซิงก์ล่าสุด
- ไม่มีการเรียก LINE API และไม่ใช้ RabbitMQ หรือ Redis
- ค่า attribute ที่ซิงก์เข้ามาถูกนำไปใช้ต่อโดย Trigger Engine และตัวกรอง audience ของ cms-api