Skip to main content

การซิงก์ข้อมูลจาก BigQuery เข้าระบบ

ภาพรวม

ลูกค้าองค์กรจำนวนมากเก็บข้อมูลลูกค้าไว้ใน Google BigQuery เช่น ยอดซื้อสะสม กลุ่มลูกค้า หรือคะแนน RFM ระบบจึงเปิดให้เชื่อมต่อ BigQuery ขององค์กรเข้ามา แล้วซิงก์ค่าเหล่านั้นลง line_user.custom_attribute เป็นประจำ ทำให้สามารถแบ่งกลุ่มผู้ใช้และตั้ง trigger จากข้อมูลฝั่ง data warehouse ได้

จุดเด่นของงานนี้คือ ตารางเวลาการทำงานถูกกำหนดโดยข้อมูลในฐานข้อมูล ไม่ใช่โดยโค้ด แต่ละแถวใน bigquery_sync_config มี cron expression ของตัวเอง และใช้ credential ของ tenant ตัวเอง

Business Flow

  1. Profile bigquery-sync เริ่มทำงาน โดยใช้เพียง PostgreSQL อย่างเดียว ไม่ต้องพึ่ง RabbitMQ หรือ Redis
  2. โหลดทุกแถวใน bigquery_sync_config ที่เปิดใช้งานอยู่ แล้วลงทะเบียน cron job หนึ่งตัวต่อหนึ่งแถว ในชื่อรูปแบบ bigquery-sync-<id> โดยใช้ค่า sync_cron ของแถวนั้น หรือใช้ค่าเริ่มต้น 0 6 * * * ซึ่งหมายถึงตี 6 ของทุกวัน
  3. เมื่อถึงเวลาทำงาน ระบบจะ
    • อ่าน config ใหม่จากฐานข้อมูลอีกครั้ง เผื่อมีการแก้ไขหลังลงทะเบียน cron ไปแล้ว
    • สร้าง BigQuery client เฉพาะ tenant จาก project_id และ credentials_json
    • Query ตารางสรุปตามที่ config ระบุไว้
  4. จับคู่และเขียนข้อมูล
    • ค้นหา line_user จากค่าในคอลัมน์ที่กำหนดใน user_id_column
    • แปลงชื่อคอลัมน์ของ BigQuery ไปเป็นคีย์ของ custom_attribute ตามที่ระบุใน field_mappings
    • เขียนข้อมูลแบบ jsonb merge ทีละชุด ชุดละ 100 แถว
    • แถวที่เกิดข้อผิดพลาดจะถูกข้ามและนับเก็บไว้ โดยไม่ทำให้ทั้งงานล้มเหลว
  5. เขียนค่า last_sync_at และ last_sync_result ซึ่งบรรจุจำนวนแถวที่อัปเดตสำเร็จและจำนวนข้อผิดพลาด กลับลงแถว config เพื่อให้ CMS แสดงสถานะการซิงก์ล่าสุดได้
  6. ตอนปิดระบบ ตัว scheduler จะรอให้ job ที่กำลังทำงานอยู่เสร็จก่อนจึงหยุด

ไฟล์และฟังก์ชันหลัก

ไฟล์หน้าที่
internal/bigquerysync/service.goService.Start และ Service.Stop ควบคุมวงจรชีวิตของ cron scheduler ส่วนการทำงานจริงอยู่ที่ loadAndRegisterConfigs, registerCronForConfig, syncForConfig, createBigQueryClient, querySummaryTable, batchUpdateUsers, updateUserAttribute, updateSyncResult พร้อมค่าคงที่ขนาด batch 100 แถว และ cron เริ่มต้น 0 6 * * *
internal/bigqueryx/bigqueryx.goWrapper ของ Google BigQuery client
cmd/worker/main.gorunBigquerySync — จุดเริ่มของ 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