Bổ sung nguồn LaoID để tạo segment theo kiến trúc v2. Dữ liệu lấy từ LaoID đồng thời đi vào bể khách hàng chung của SP; Excel, OCS và campaign event dùng cùng quy tắc nhận diện để không tạo nhiều hồ sơ cho cùng khách hàng.
Task gồm tạo segment LaoID, hồ sơ/định danh theo SP, chuẩn bị thông tin cho campaign event và nối sự kiện về customer để lọc tệp theo hành vi. Đây là các phần cùng một task. Campaign event là luồng chính; không yêu cầu sử dụng PostHog. Không thay đăng nhập/quyền LaoCRM hoặc triển khai phần base tổng quát.
/segment/data-sources, SegmentSourceEnum, SegmentVersionSourceTypeEnum và SegmentBuildJobTypeEnum.POST /segment hiện có với type=QUERY, source=LAOID và data/queryDefinition cho điều kiện lọc. SP lấy từ context đã xác thực, không tin mã SP tùy ý trong body.{
"name": "Khách hàng LaoID",
"type": "QUERY",
"source": "LAOID",
"data": [
{ "attribute": "status", "operator": "EQUALS", "value": "1", "relation": "AND" }
]
}
Trường/operator trong ví dụ chỉ được dùng nếu có trong metadata đã được cấp quyền. Dùng lại metadata user_lao_id_field_config/API field hiện có, bổ sung phân quyền và whitelist; client không gửi table/column/join hoặc raw SQL.
segment-version-import hiện có; thêm nhánh LaoID trong segment-import-worker với DataSource read-only có pool/timeout riêng, SQL bind giá trị và cursor/batch có giới hạn.Sửa đường query trong phạm vi task, không chỉ tăng DEFAULT_LIMIT. Source hiện SegmentService.queryUserInfo() gọi JDBC và tạo List cho toàn kết quả trả về; CustomerQueryBuilder dùng SELECT * và ghép giá trị vào SQL. Không dùng đường này để build tập LaoID lớn trong API.
Luồng đọc bắt buộc:
Đọc batch đầu
→ Chuẩn hóa/resolve và ghi batch
→ Sink xác nhận, lưu khóa cuối vào checkpoint
→ Đọc tiếp từ khóa cuối
→ Hết dữ liệu, kiểm tra và publish version
SELECT id, primary_phone, primary_email, first_name
FROM user_basic_info
WHERE id > :lastId
ORDER BY id ASC
LIMIT :batchSize
Batch đầu bỏ điều kiện lastId, :batchSize mặc định 3000. Ví dụ này cần tên; không tự lấy first_name cho mọi job. Chỉ SELECT field thực sự cần/được phép; không SELECT * hoặc giữ một LIMIT tổng làm cắt danh sách mà vẫn báo hoàn tất. Cursor/queryDefinitionChecksum/attempt không được đổi giữa job.
Tham chiếu kỹ thuật: MariaDB keyset pagination, EXPLAIN.
Backend lập danh sách cột và mapping alias trước khi chạy job/enrichment, dựa trên metadata được cấp quyền. Không lấy mọi field LaoID vì có trong schema hoặc biến có tên giống cột.
Cột cursor/định danh cần dùng
+ Cột cần lọc tại worker
+ Trường hồ sơ được chọn để cập nhật
+ Biến template/action thực sự lấy từ LaoID
= Danh sách SELECT
| Nhóm / biến | Cột nguồn / xử lý |
|---|---|
| Cursor và định danh nguồn | id bắt buộc; giữ kiểu text/ID nguồn khi mapping, không thay bằng customerId |
| Hợp nhất contact | primary_phone và primary_email là bộ nền khi nguồn/policy cho phép đối chiếu cả hai. Email có thể cần cho identity dù chỉ gửi SMS; null contact không loại user email-only/ID-only |
{{first_name}}, {{last_name}}, {{username}} |
first_name → firstName, last_name → lastName, username → username; chỉ lấy khi template/field hồ sơ đã chọn cần |
{{primary_phone}}, {{primary_email}} |
primary_phone → primaryPhone, primary_email → primaryEmail; alias phone/email nếu hỗ trợ tương thích phải map rõ về cùng cột, không query hai lần |
{{avatar}} |
avatar → avatar chỉ khi được phép/cần dùng; giữ quy tắc ghép image endpoint hiện có |
| Status, country/province/address/language… | Lấy nếu lọc tại worker, template hoặc hồ sơ đã chọn cần; filter hoàn toàn trong SQL không tự yêu cầu SELECT cột đó nếu không dùng tiếp |
| Biến event/campaign như amount, package, renew_percentage | Dùng request/config khi đã có và được phép; không tự coi alias là cột LaoID |
{{track.key}} |
Do tracking tạo, không query LaoID |
Đọc biến từ SMS content; mail title/content; WhatsApp content/templateParameters; action ObjectOrRef như amount/period/tradeMark; condition parameter và tracking originalUrl có biến. Dùng cùng parser/ref semantics với renderer, không chỉ regex một trường SMS. Khai báo source của alias (LaoID/request/profile/campaign/tracking/derived), mapping và dataType; identity/phone/email/ID số dạng chuỗi không tự infer thành NUMBER rồi mất số 0 đầu.
Source renderer SMS/mail dùng event.params: dữ liệu đã lấy phải map vào params đúng alias rồi chốt; không chỉ gán firstName vào UserInfoModel và kỳ vọng template tự thấy. Field không có trong model, ví dụ country/province nếu cần, phải map vào attributes/params có whitelist. Không chỉ thêm cột SELECT rồi bỏ dữ liệu khi mapping.
Sửa các chỗ hiện có:
{{first_name}} hoặc tự trả chuỗi rỗng cho trường bắt buộc.Khi tạo segment chưa gắn campaign, selected fields đến từ danh sách hồ sơ/template được chọn và field lọc. Khi campaign dùng version, kiểm required aliases/contact đã có trong payload; nếu thiếu thì lấy từ dữ liệu địa phương hoặc tạo version/snapshot bổ sung trước chạy theo quyền. Không âm thầm query nguồn cho từng người trong worker hoặc cam kết template nào cũng đủ dữ liệu từ version tối thiểu.
Lưu selectedSourceFields, requiredParamAliases, fieldMappingVersion và queryPlan vào source snapshot/job; kiểm quyền nguồn/field trước SELECT, không để client truyền nguyên tên cột. Trường bắt buộc cho xử lý nội bộ vẫn cần quyền service đối với nguồn; quyền trả đầy đủ cho người dùng được kiểm riêng ở phần bảo mật. Không SELECT password/hash, token, toàn danh sách phones/emails hoặc attributes lớn nếu không cần/không được cấp quyền.
Các nhánh cùng dùng resolver nhưng không mặc định dòng import là hành vi. Payload shadow giữ thông tin đã chốt; membership giữ danh sách thuộc version; profile giữ thông tin hiện tại.
Mỗi SP có customer_profile_<storageId> và customer_identity_<storageId>, do backend tạo theo schema/index chuẩn. Registry ánh xạ SP đã xác thực sang tên collection ổn định; frontend không được chọn tên collection. Các bảng hành vi/segment ClickHouse dùng chung, mọi query có spCode.
Một customerId nội bộ có nhiều identity: phone, email hoặc ID nguồn có namespace. ID LaoID=123 không tự khớp ID OCS=123. User chỉ có email vẫn hợp lệ; tên/địa chỉ không dùng để gộp user.
Giữ xác thực endpoint/campaign/service hiện tại. Kết quả lấy thêm thông tin từ LaoID hoặc nguồn khác phải được chuẩn hóa và đưa qua resolver chung để cập nhật hồ sơ SP, không chỉ dựng user tạm cho request.
Luồng: xác thực → chuẩn hóa → bổ sung nguồn nếu cần → đối chiếu toàn identity → cập nhật hồ sơ → chốt customerId/profileVersion/params → dispatcher/worker. Tách FOUND/NOT_FOUND/ERROR của nguồn; không phone-first bỏ qua ID mâu thuẫn. Cách tối ưu query nguồn được cấu hình qua adapter, không gọi nguồn lại trong từng action.
Query bổ sung thông tin cho event cũng nằm trong task: ưu tiên request/local profile đủ field và độ mới cho phép; cần nguồn thì lookup bằng phone/email/userId có index, số kết quả hữu hạn, pool/concurrency/timeout riêng và cache chống nhiều request cùng lookup. Không SELECT toàn nguồn hoặc áp hàm normalize lên cả cột rồi scan theo từng event; format query phải khớp cách nguồn lưu. Không tìm thấy khác source ERROR; cơ chế cache không mặc định giữ mọi user hoặc vô thời hạn. Giữ hợp đồng phản hồi hiện có, không tự đổi sang async nếu caller chưa chấp thuận.
customerId là trường riêng, không thay UserInfoModel.id đang dùng làm ID LaoID/provider. Event field hợp lệ kết hợp profile SP/default để render; thiếu contact bắt buộc của kênh thì lỗi rõ. Không đổi params của command cũ khi profile cập nhật hoặc Kafka phát lại.
Lưu sự kiện nguồn gắn customerId ở ClickHouse, có SP/source/sourceEventId/thời gian và properties được phép. Một sự kiện kích hoạt nhiều campaign chỉ tính hành vi nguồn một lần; execution/result gửi tin riêng, không tự coi là purchase. Chưa xác định customer thì giữ identity context/unresolved để xử lý, không đoán user.
Projection hồ sơ và số liệu customer phục vụ lọc tệp, ví dụ đã đăng ký gói nhưng chưa gia hạn trong cửa sổ thời gian. Builder dùng condition v2, scope SP và dataAsOf, tạo membership/payload mới rồi publish; không thay run đang chạy. Aggregate xử lý event trùng/revision, không SUM lại các lần replay/rebuild. Không thêm journey/ML/lookalike trong task này.
| Vị trí | Nội dung cần triển khai |
|---|---|
segment / segment_version |
Source LAOID, điều kiện và snapshot theo SP hiện có; thêm selectedSourceFields/requiredParamAliases/fieldMappingVersion/queryPlan; consistencyMode/sourceSnapshotRef/readStartedAt/readCompletedAt nếu cần. Giữ version và activation v2 |
segment_build_job |
Nhánh LAOID; thêm sourceCursor/cursorField hoặc cursor tuple/queryDefinitionChecksum/checkpointAttemptId/checkpointCommittedAt/duplicateCount/conflictCount, dùng lease/count/error hiện có. Checkpoint chỉ sau batch ghi xác nhận; snapshot/boundary nằm trong source snapshot |
| Catalog/field config | Seed LAOID; alias/sourceColumn/fieldName hoặc attribute mapping/dataType, nguồn alias, field/operator được phép, country/namespace và quyền filter/template/view/export. Không credential/raw SQL trong metadata trả client |
customer_storage_registry — mới |
spCode, storageId, collectionNames, status, schema/indexVersion; khởi tạo đủ và READY trước ghi |
| Profile collection của SP | customerId/spCode, canonical/mergedInto, tên/attributes/status/profileVersion, fieldMeta nguồn/thời điểm/priority, contact policy và timestamps. Không toàn event history |
| Identity collection của SP | customerId, type/namespace, normalized hash/encrypted value, normalizationVersion, source/assurance, status/validFrom/validTo/version. Unique active type+namespace+hash trong SP; index customerId |
| Conflict/merge audit — metadata chung | spCode, candidate/target customerIds, identity hashes, source/policy, reason/status, actor/bằng chứng/timestamps; chỉ case cần xử lý |
| Event/request/member DTO | Thêm email/typed identifiers/sourceEventId nếu thiếu; customerId/profileVersion riêng; giữ alias/ID nguồn và response hiện tại |
| ClickHouse profile/identity projection | Dùng bảng chung theo SP; version và binding time để lọc/resolve, thêm field whitelist cần thiết. Source không chạy hai sink trùng |
| ClickHouse behavior/aggregate | Event nguồn customerId nullable/identity context, kind/revision/time; tổng hợp customer/time có generation/dataAsOf và publish hoàn tất |
| Payload/membership hiện có | customerId/profileVersion và dữ liệu chốt; member_id vẫn là mã payload, không customerId. Dedup customer/recipient trong version/attempt |
Các module: API/metadata/event integration trong laoads; adapter/import trong segment-import-worker; model/resolver/repository cụ thể dùng chung trong common; projection/behavior và query tệp trong analytics/audience v2. App-identify hiện có phải dùng cùng router/resolver khi liên quan, không tiếp tục tạo kho hồ sơ khác.
Quyền truy cập SP và quyền xem dữ liệu là hai điều kiện riêng. Được thao tác SP A không mặc định được xem đầy đủ phone/email của user trong A.
| Thao tác | Kiểm tra bắt buộc |
|---|---|
| Tạo/sửa/refresh segment LaoID | Quyền thao tác segment của SP + quyền sử dụng nguồn/tập dữ liệu LaoID và field lọc |
| Danh sách/chi tiết segment, version, job, số lượng | Quyền xem segment và owner SP của đối tượng, kể cả truy bằng ID |
| Xem/search/preview thành viên | Quyền xem thành viên; field được phép, che phone/email/thông tin cá nhân khi không có quyền xem đầy đủ |
| Xuất danh sách/file lỗi | Quyền export riêng và quyền field tương ứng; không có quyền thì chặn, không dựa vào quyền xem màn hình |
| Xem hồ sơ/identity hoặc xử lý conflict/merge | Quyền dữ liệu khách hàng trong SP và quyền chức năng tương ứng; merge có audit, không tự cấp từ VIEW |
| Chạy campaign | Quyền sử dụng segment của SP; worker đọc dữ liệu cần gửi bằng service permission, không trả contact đầy đủ cho caller không có quyền |
Dùng quyền LaoCRM/v1 hiện có hoặc bổ sung quyền chức năng tương ứng; không tự cấp quyền mọi dữ liệu cho tất cả user cùng SP. Chốt mapping quyền trước bật API, không chỉ isAuthenticated().
| Trường hợp | Kết quả |
|---|---|
| Tạo QUERY/LAOID hợp lệ | Trả segment/version/job, worker build async; publish xong chọn chạy campaign |
| Field/operator sai hoặc ngoài quyền | Từ chối trước query nguồn; không thực thi SQL tùy ý |
| LaoID không có dòng phù hợp | Segment rỗng đúng kết quả; phân biệt nguồn lỗi/timeout |
| Import khoảng 10 triệu dòng, dừng/phát lại | Keyset/batch có index và giới hạn, checkpoint sau ghi xác nhận, không OFFSET lớn/List toàn dữ liệu; restart không mất/trùng dòng hoặc publish phần dở |
| Filter rộng/query thiếu index/downstream chậm | Kiểm query plan/budget, giới hạn hoặc chờ, không tăng concurrency để làm quá tải nguồn; worker ngừng đọc khi sink nghẽn |
| Preview/count/lookup event trên nguồn lớn | Query riêng có giới hạn/index/cache/timeout và đúng quyền; không scan toàn nguồn mỗi thao tác hoặc biến timeout thành 0/không có user |
| Batch 3000, nhiều batch/job hoặc identity keys toàn nguồn | Giới hạn in-flight/heap và trường/attributes, giải phóng batch; không HashSet tăng theo cả triệu user |
| Template chỉ cần tên, amount đã có trong event | SELECT bộ identity + tên + field thực sự cần, không query amount/avatar/cột khác tự động |
| Alias không mapping/không quyền/cột chưa SELECT | Từ chối hoặc fallback đã chốt trước gửi; không SQL từ alias, mapper không đọc cột thiếu và không xóa field cũ |
| Campaign mới cần field chưa có trong version | Kiểm tra trước chạy, chuẩn bị snapshot/version bổ sung theo quyền; không gọi nguồn lại cho từng action |
| Nhiều format phone, cùng email/ID đã liên kết | Normalize và resolve về một profile khi không conflict |
| Email-only hoặc ID-only | Hồ sơ hợp lệ theo policy; chỉ action có đủ contact mới chạy |
| Phone trỏ A, email/ID trỏ B | Conflict; xác minh/merge có căn cứ, không chọn một khóa rồi bỏ qua khóa khác |
| Hai nguồn/request đồng thời cùng identity | Unique/re-resolve, không hai profile hoặc orphan identity |
| Query identity/enrichment lỗi | Không coi user mới hoặc xóa dữ liệu cũ; fallback/validation theo field cần thiết |
| Contact đổi chủ/merge/profile đổi | Binding/audit đúng thời điểm, không gán lại hành vi sai hoặc thay command/run đã chốt |
| Event trùng hoặc kích hoạt nhiều campaign | Một hành vi nguồn, execution riêng; không thu/chạy lại vì replay |
| User thuộc SP A truy object/file/cache SP B | Từ chối, không trả dữ liệu hoặc metadata nhạy cảm |
| Có quyền xem segment, thiếu quyền dữ liệu/export | Chỉ trả trường được phép/masked; chặn export hoặc xem đầy đủ |
Phần thao tác dữ liệu phải dùng index/query scope, batch và giới hạn nguồn theo SP/toàn hệ thống. Cấu hình collection riêng không tự bảo đảm tải 10 triệu customer/SP; việc chọn shard/hạ tầng không tự thực hiện trong task code. Áp rule timeout tiền/action và pipeline ACK của các task chung.