Bỏ qua để đến nội dung

Cloudera Data Engineering: Apache Spark, điều phối Airflow trên Iceberg

Dữ liệu hiếm khi dùng được ngay khi vừa lấy về. Nó đến từ nhiều hệ thống với định dạng khác nhau, thiếu trường, trùng bản ghi, sai kiểu, và mỗi nguồn gọi cùng một khách hàng bằng một cái tên khác. Khoảng cách giữa “đã có dữ liệu” và “dữ liệu dùng được để ra quyết định” chính là phần việc của Data Engineering.

Cloudera Data Engineering là dịch vụ để đội kỹ thuật dữ liệu dựng và điều phối luồng xử lý trên nền Apache Spark, điều phối bằng Apache Airflow, làm việc trực tiếp trên bảng Apache Iceberg của lakehouse. Toàn bộ vòng đời nằm trong một chỗ: viết logic biến đổi, thử tương tác, chạy theo lịch có phụ thuộc, quản lý phiên bản dữ liệu, và theo dõi khi có sự cố.

Nhìn từ trên xuống, Data Engineering đứng giữa hai đầu — dữ liệu vào từ mọi nguồn ở bên trái, kết quả phục vụ các tình huống nghiệp vụ ở bên phải — và làm ba nhóm việc: Biến đổi, Điều phối, Tối ưu. Tất cả đặt trên nền Open Data Lakehouse (Apache Iceberg) và mang theo được giữa trung tâm dữ liệu và đám mây.

Ba trụ Cloudera Data Engineering: Biến đổi, Điều phối, Tối ưu

Sơ đồ này bổ sung vài chi tiết đáng lưu ý mà các phần dưới sẽ nói kỹ hơn: Spark chạy đóng gói trong container (mỗi luồng có môi trường riêng, không giẫm chân nhau), hỗ trợ đa ngôn ngữ (SQL, Python, Java — đội không phải viết lại theo một ngôn ngữ áp đặt), API tự động hóa pipeline cùng tích hợp CI/CD và Git (luồng dữ liệu đi qua đúng quy trình phát hành phần mềm như mã nguồn khác), và chẩn đoán thời gian thực bên cạnh gỡ lỗi trực quan. Phần còn lại của trang đi vào từng trụ.

flowchart LR
  SRC["Nguồn dữ liệu<br/>core banking · ERP · log · tệp"] --> ING["Nạp dữ liệu<br/>theo lô hoặc theo luồng"]
  ING --> RAW["Bảng Iceberg thô"]
  RAW --> SPARK["Apache Spark<br/>làm sạch · biến đổi · đối chiếu"]
  SPARK --> CUR["Bảng Iceberg chuẩn hóa"]
  CUR --> SERVE["Lớp phục vụ<br/>báo cáo · dashboard · mô hình AI"]
  ORCH["Apache Airflow<br/>DAG · phụ thuộc · thử lại · SLA"] -.-> SPARK
  MON["Theo dõi<br/>trạng thái · khối lượng · tài nguyên"] -.-> SPARK
  SDX["SDX — chính sách & truy vết tự động"] -.-> SPARK

Phần lớn thời gian của một dự án dữ liệu không nằm ở mô hình hay dashboard, mà nằm ở đoạn đường giữa nguồn và đích. Data Engineering lo bốn nhóm việc:

Nhóm việcCụ thể là gì
Làm sạchBỏ bản ghi trùng, chuẩn kiểu dữ liệu, xử lý trường thiếu, bắt giá trị bất thường trước khi nó vào báo cáo
Biến đổiGhép nhiều nguồn, tính chỉ tiêu nghiệp vụ, dựng bảng tổng hợp cho lớp phục vụ
Điều phốiChạy đúng thứ tự phụ thuộc, tự thử lại khi nguồn về trễ, báo khi trễ cam kết thời gian
Theo dõiBiết job nào chạy bao lâu, hỏng ở đâu, ảnh hưởng bảng nào — trước khi người dùng phát hiện ra

Điểm khác biệt của cách làm trên lakehouse: đầu ra không phải một bản trích xuất gửi đi đâu đó, mà chính là bảng Iceberg mà báo cáo và mô hình AI (Artificial Intelligence — trí tuệ nhân tạo) sẽ đọc. Một bản dữ liệu, một nơi để sửa khi sai.

Cloudera đóng khung phần việc này quanh ba thách thức của một đội dữ liệu, và cách đóng khung đó khớp với thực tế các dự án tại Việt Nam:

Thách thứcNền tảng trả lời bằng gì
Gom được nhiều nguồn khác nhauSpark xử lý cả dữ liệu theo lô lẫn theo luồng trong một khung; Iceberg làm định dạng bảng chung; lakehouse mở quản lý thống nhất qua môi trường hybrid
Điều phối luồng ở quy mô lớnAirflow điều phối theo phụ thuộc, kèm tích hợp API tới công cụ ngoài; kiểm soát truy cập chi tiết và truy vết nguồn gốc tự động
Kiểm soát chi phí trên các khối lượng công việcBộ công cụ DataOps, lớp quan sát và quản trị tài chính ở mức khối lượng công việc

Apache Spark là bộ máy xử lý dữ liệu phân tán — chia khối lượng công việc ra nhiều máy để xử lý dữ liệu lớn trong thời gian chấp nhận được. Theo Cloudera công bố, Spark ở đây là một khung xử lý thống nhất cho cả dữ liệu theo lô (batch) và theo luồng (streaming) — nghĩa là không cần hai bộ máy và hai bộ kỹ năng cho hai kiểu dữ liệu. Đây là chi tiết đáng cân nhắc khi lập kế hoạch nhân sự cho đội dữ liệu, không chỉ là một dòng thông số.

Trong Cloudera, Spark làm việc trực tiếp trên bảng Iceberg, và sự kết hợp đó mang lại vài thứ mà cách làm cũ trên hồ dữ liệu không có:

Ghi có giao dịch. Iceberg đảm bảo ACID (Atomicity, Consistency, Isolation, Durability — nguyên tử, nhất quán, cô lập, bền vững): job Spark chết giữa chừng thì bảng không để lại dữ liệu dở dang. Không cần dọn tay, không cần bảng tạm.

Ghi không chặn đọc. Job nạp liệu ban đêm chạy song song với người đang truy vấn. Người đọc thấy một phiên bản nhất quán, không thấy trạng thái đang ghi giữa chừng.

Cập nhật và xóa theo dòng. Sửa một tập bản ghi cụ thể mà không phải ghi lại cả phân vùng — rất cần cho luồng chỉnh sửa hồi tố, hoặc khi phải xóa dữ liệu của một cá nhân theo quy định bảo vệ dữ liệu cá nhân.

Đổi schema không gãy luồng. Nguồn thêm trường mới thì thêm cột vào bảng đích; những job và truy vấn viết từ trước vẫn chạy.

Một luồng dữ liệu bền hiếm khi biến đổi thẳng từ nguồn ra bảng báo cáo. Cách làm phổ biến là chia bảng thành ba tầng có vai trò rõ ràng:

TầngChứa gìVì sao cần
ThôDữ liệu giữ gần nguyên trạng như lúc nhận từ nguồnKhi phát hiện logic biến đổi sai, có chỗ để chạy lại — không phải đi xin nguồn cấp lại dữ liệu cũ
Chuẩn hóaĐã làm sạch, khử trùng lặp, thống nhất mã và kiểu dữ liệuĐây là bản dữ liệu đáng tin để mọi đội cùng dùng, thay vì mỗi đội tự làm sạch một kiểu
Phục vụBảng tổng hợp, bảng đặc trưng dựng theo nhu cầu báo cáo và mô hìnhTách nhu cầu hiển thị khỏi nhu cầu lưu trữ; đổi báo cáo không phải đụng tầng dưới

Ranh giới này trả lời được câu hỏi hay gặp nhất khi số liệu sai: sai từ nguồn, sai ở bước biến đổi, hay sai ở cách tổng hợp? Có ba tầng thì khoanh vùng trong vài phút. Gộp tất cả vào một bước thì chỉ còn cách đọc lại toàn bộ mã.

Một tiêu chí ít được nói tới nhưng quyết định việc bạn đặt ranh giới đúng chỗ: ba tầng có tốc độ thay đổi khác nhau. Tầng thô gần như không đổi — đổi nghĩa là mất bản gốc. Tầng chuẩn hóa đổi chậm và có kiểm soát, vì nó là hợp đồng dùng chung giữa các đội. Tầng phục vụ đổi thường xuyên, và đó là chuyện bình thường. Nếu báo cáo đổi mà bạn phải sửa tầng chuẩn hóa, ranh giới đã bị đặt sai chỗ.

Vì sao đặt lịch rời rạc là một cái bẫy

Phần tiêu đề “Vì sao đặt lịch rời rạc là một cái bẫy”

Cách bắt đầu tự nhiên nhất của mọi đội: mỗi job một lịch chạy. Nạp lúc 1 giờ, chuẩn hóa lúc 2 giờ, dựng bảng phục vụ lúc 3 giờ. Chừa mỗi bước một tiếng cho chắc.

Cách này chạy tốt — cho tới ngày một nguồn về trễ mười phút.

Lúc đó bước nạp chưa xong, nhưng đồng hồ điểm 2 giờ và bước chuẩn hóa vẫn cứ chạy, trên dữ liệu thiếu. Nó thành công. Bước dựng bảng phục vụ lúc 3 giờ cũng thành công. Sáng ra, báo cáo có số, số thiếu một nguồn, và cả ba job đều xanh. Không có gì trong hệ thống nói rằng vừa có chuyện gì đó sai.

Nguyên nhân gốc gói trong một câu: giờ chạy không phải là quan hệ phụ thuộc. Một tiếng chừa ra là một lời phỏng đoán về thời gian, và mọi phỏng đoán về thời gian đều sai vào đúng ngày bận nhất — tức là đúng ngày con số bị soi kỹ nhất.

Airflow: khai báo phụ thuộc, không khai báo giờ

Phần tiêu đề “Airflow: khai báo phụ thuộc, không khai báo giờ”

Cloudera Data Engineering dùng Apache Airflow làm bộ máy điều phối (theo Cloudera công bố). Airflow đảo ngược cách đặt vấn đề: bạn không nói “chạy lúc 2 giờ”, bạn nói “chạy sau khi bước nạp xong”.

Cả luồng được mô tả thành một DAG (Directed Acyclic Graph — đồ thị có hướng không chu trình): các bước là nút, quan hệ phụ thuộc là cạnh, không có chu trình vì dữ liệu không chảy ngược. Bộ điều phối đọc đồ thị và tự biết bước nào chạy được, bước nào phải chờ.

flowchart LR
  subgraph CRON["Đặt lịch rời rạc"]
    direction TB
    A1["1:00 nạp"] -.->|"chỉ là phỏng đoán<br/>về thời gian"| A2["2:00 chuẩn hóa"]
    A2 -.-> A3["3:00 bảng phục vụ"]
    A3 --> A4["Nguồn trễ 10 phút<br/>⇒ 3 job xanh, số sai"]
  end
  subgraph DAG["DAG trên Airflow"]
    direction TB
    B1["nạp"] -->|"phụ thuộc thật"| B2["chuẩn hóa"]
    B2 --> B3["đối chiếu"]
    B3 --> B4["bảng phục vụ"]
    B4 --> B5["Nguồn trễ 10 phút<br/>⇒ cả chuỗi dịch 10 phút, số đúng"]
  end

Nguồn về trễ mười phút? Bước chuẩn hóa chờ mười phút rồi mới chạy — trên dữ liệu đủ. Nguồn không về? Chuỗi dừng ở đúng chỗ đó và báo, thay vì chạy tiếp trên dữ liệu thiếu.

Thử lại có chính sách. Lỗi tạm thời — nghẽn mạng, nguồn bận — không cần con người. Đặt số lần thử lại và khoảng cách giữa các lần; chỉ khi hết số lần mới gọi người.

Cam kết thời gian (SLA — Service Level Agreement). Khai báo “chuỗi này phải xong trước 6 giờ sáng”. Chưa xong thì cảnh báo phát đi lúc 6 giờ — không phải lúc 9 giờ khi người dùng phát hiện ra. Đây là khác biệt giữa việc bạn báo cho nghiệp vụ và nghiệp vụ báo cho bạn.

Chạy bù (backfill). Sửa xong logic thì cho chạy lại đúng khoảng thời gian bị ảnh hưởng, theo đúng thứ tự phụ thuộc, không phải viết một script riêng cho mỗi lần sự cố.

Tích hợp ra ngoài nền tảng. Theo Cloudera công bố, Airflow ở đây có sẵn tích hợp API (Application Programming Interface — giao diện lập trình ứng dụng) với các công cụ dữ liệu bên ngoài. Điều này quan trọng hơn vẻ ngoài kỹ thuật của nó: luồng thật ở một ngân hàng hiếm khi chỉ gồm các bước bên trong nền tảng — nó còn phải kích hoạt một công việc ở hệ thống khác, chờ tín hiệu từ hệ thống lõi, hoặc báo cho một công cụ nghiệp vụ khi dữ liệu sẵn sàng. Điều phối được cả những bước nằm ngoài nghĩa là cả chuỗi có một chỗ để nhìn, thay vì một nửa nằm trong bộ điều phối và một nửa nằm trong trí nhớ của người vận hành.

Theo Cloudera công bố, phần điều phối, quản lý và đặt lịch ở đây là tự phục vụ: đội dữ liệu tự dựng và tự chạy luồng của mình, không phải mở phiếu yêu cầu gửi đội nền tảng.

“Tự phục vụ” thường bị hiểu là “buông”. Không phải — và chỗ này đáng nói rõ vì nó quyết định bạn mở rộng nền tảng được tới đâu. Quyền xem dữ liệu vẫn nằm ở lớp quản trị dùng chung SDX: đội nào tự dựng luồng gì thì luồng đó vẫn chỉ đọc được đúng phần dữ liệu mà đội đó được phép đọc.

Tự phục vụ ở tầng điều phối, kiểm soát ở tầng dữ liệu. Hai thứ này nằm ở hai lớp khác nhau, và tách được chúng ra chính là điều cho phép thêm đội dùng nền tảng mà không phải nới lỏng bên nào.

Đây là chỗ bảng Iceberg thay đổi hẳn cách xử lý sự cố. Mỗi lần job ghi vào bảng, Iceberg tạo một ảnh chụp (snapshot) — một phiên bản bất biến của bảng tại thời điểm đó.

  • Job nạp sai dữ liệu? Trả bảng về ảnh chụp liền trước. Vài giây, không cần phục hồi từ bản sao lưu, không cần dừng hệ thống.
  • Số liệu tuần này lệch tuần trước? Đọc lại bảng ở cả hai thời điểm rồi so — thay vì suy đoán.
  • Kiểm toán hỏi số liệu cuối quý? Truy vấn đúng ảnh chụp cuối quý đó, không tranh cãi.

Đặt cạnh nhau hai quy trình xử lý cùng một sự cố:

Cách cũVới ảnh chụp
Phát hiện job nạp saiThường sau khi báo cáo đã phát hànhBắt được ở bước đối chiếu, trước khi ghi
Khôi phụcTìm bản sao lưu, phục hồi, chạy lạiQuay lui về ảnh chụp liền trước
Thời gianTính bằng giờ, đôi khi bằng ngàyTính bằng giây
Hệ thống có phải dừngThường là cóKhông
Giải trình sau đóViết lại từ trí nhớẢnh chụp là bản ghi sẵn có

Theo dõi: job xanh không có nghĩa là số đúng

Phần tiêu đề “Theo dõi: job xanh không có nghĩa là số đúng”

Luồng dữ liệu hỏng theo nhiều kiểu, và kiểu nguy hiểm nhất là kiểu không báo lỗi: job vẫn xanh, nhưng nguồn hôm nay chỉ về một nửa số bản ghi. Vì vậy theo dõi cần đi xa hơn trạng thái thành công hay thất bại. Bốn tín hiệu đáng nhìn cho mỗi luồng:

Trạng thái — job chạy xong hay hỏng, và hỏng ở bước nào trong chuỗi.

Thời gian chạy so với chính nó những ngày trước — không so với một ngưỡng cố định. Một job thường mất 20 phút mà hôm nay xong sau 3 phút cũng đáng ngờ ngang với job chạy 2 tiếng. Nhiều đội chỉ đặt ngưỡng cho “chậm hơn X” mà quên chiều còn lại — trong khi nhanh bất thường gần như luôn có nghĩa là thiếu dữ liệu.

Khối lượng vào ra ở từng bước — đây là tín hiệu duy nhất bắt được kiểu hỏng thầm lặng. Trạng thái job, theo thiết kế, không bao giờ nói ra chuyện này.

Tài nguyên tiêu tốn — để biết chi phí đang đi về đâu trước khi hóa đơn cuối tháng nói cho biết.

Cloudera Data Engineering giữ nhật ký chạy và số liệu của từng job để đội vận hành khoanh vùng nhanh; phần giám sát ở mức toàn nền tảng nằm ở lớp Observability.

Đối chiếu là một bước của luồng, không phải việc làm sau

Phần tiêu đề “Đối chiếu là một bước của luồng, không phải việc làm sau”

Sai lầm phổ biến: coi kiểm tra chất lượng là việc rà soát hằng tuần. Lúc đó dữ liệu sai đã nằm trong bảng phục vụ bảy ngày và đã đi vào ba báo cáo cùng một mô hình.

Cách bền hơn: đặt bước đối chiếu vào giữa luồng, trước khi ghi vào tầng phục vụ. Lệch thì luồng dừng và cảnh báo; bảng phục vụ giữ nguyên số của hôm qua — và số của hôm qua thì đúng.

Phép đối chiếuBắt được gì
Đếm số bản ghi so với nguồnMất dữ liệu, nạp trùng
Tổng kiểm các cột số (control total)Sai lệch giá trị, sai đơn vị, lỗi ép kiểu
So mẫu ngẫu nhiên với nguồnSai logic biến đổi mà tổng vẫn khớp
Kiểm mốc biên: cuối kỳ, dữ liệu về trễ, múi giờKiểu lỗi chỉ xuất hiện đúng ngày quan trọng nhất

Đây là một đánh đổi có chủ ý, và nên nói thẳng: nó có nghĩa là có ngày báo cáo sẽ trễ. Đó vẫn là cuộc đổi chác hời. Một báo cáo trễ là một sự bất tiện — người ta biết nó trễ. Một báo cáo sai là một quyết định sai, đưa ra bởi một người tin rằng mình đang có thông tin đúng.

DataOps: một vòng đời, từ phát triển tới vận hành

Phần tiêu đề “DataOps: một vòng đời, từ phát triển tới vận hành”

Theo Cloudera công bố, Cloudera Data Engineering đóng gói toàn bộ vòng đời từ phát triển tới triển khai, kèm giám sát vận hành, cùng bộ công cụ DataOps để rút ngắn thời gian tạo ra giá trị. Vài mảnh trong đó thay đổi công việc hằng ngày của đội:

Phiên làm việc tương tác (interactive session). Thử một phép biến đổi và thấy kết quả ngay, thay vì đóng gói một job rồi chờ nó chạy để biết mình viết sai một dấu. Vòng lặp thử — sai — sửa ngắn lại, và đó là thứ quyết định năng suất thật của một kỹ sư dữ liệu.

Kết nối từ môi trường phát triển bên ngoài (IDE — Integrated Development Environment). Kỹ sư dùng công cụ quen tay, kết nối từ xa vào Spark trên nền tảng. Điểm đáng giá không phải sự tiện: nó là mã nguồn ở lại trong kho mã của tổ chức và đi qua đúng quy trình duyệt — thay vì sống trong những notebook rời mà không ai review.

Gỡ lỗi trực quan và tinh chỉnh hiệu năng. Nhìn thấy job chậm ở bước nào, thay vì đọc nhật ký để đoán.

Quản lý và đặt lịch tự phục vụ. Đội dữ liệu tự vận hành luồng của mình, trong khuôn chính sách của nền tảng.

Trong hồ sơ dự án, nhóm năng lực này hay bị xếp vào “công cụ nội bộ” và đẩy xuống giai đoạn sau — vì nó không tạo ra một báo cáo nào mà lãnh đạo nhìn thấy. Cái giá thì xuất hiện đều đặn và âm thầm: mỗi lần đổi một luồng là một lần hồi hộp vì không có chỗ thử; mỗi sự cố là một cuộc điều tra vì không có công cụ nhìn; và mỗi kỹ sư giỏi rời đi đều mang theo một phần hệ thống, vì phần đó chưa bao giờ rời khỏi máy của họ.

Chi phí: quản trị tài chính ở mức khối lượng công việc

Phần tiêu đề “Chi phí: quản trị tài chính ở mức khối lượng công việc”

Theo Cloudera công bố, Cloudera Data Engineering có lớp quan sát và quản trị tài chính ở mức khối lượng công việc — quy được chi phí về đúng đội, đúng dự án, đúng luồng.

Việc gắn tên vào con số thay đổi hành vi nhiều hơn bất kỳ chính sách nào. Một cụm chạy suốt ngày đêm cho một job hai tiếng sẽ tồn tại chừng nào chi phí của nó còn là “chi phí hạ tầng chung”, và biến mất khá nhanh khi nó nằm trong báo cáo chi phí của một đội cụ thể — không cần ai ra lệnh.

Điều kiện để làm được: gắn chủ sở hữu cho từng khối lượng công việc trước khi bắt đầu đo. Việc này tốn vài ngày lúc đầu và gần như bất khả thi sau hai năm, khi không ai còn nhớ luồng đó dựng cho ai.

Một lưu ý dễ bị bỏ: hóa đơn hạ tầng không phải toàn bộ chi phí. Một luồng chạy chậm cũng là tiền — chỉ là tiền trả bằng thời gian của người khác: phòng tài chính chờ số để chốt, đội phân tích chờ dữ liệu để trả lời lãnh đạo. Khoản đó không nằm trên hóa đơn nào, nên nó thường thắng mọi cuộc tranh luận về tối ưu — cho tới khi có người quy nó ra tiền.

Bảo mật và truy vết kế thừa từ nền tảng

Phần tiêu đề “Bảo mật và truy vết kế thừa từ nền tảng”

Cloudera Data Engineering không tự dựng một bộ quyền riêng. Nó là một dịch vụ trên nền tảng nên kế thừa lớp quản trị dùng chung SDX: chính sách đặt một lần, áp cho mọi dịch vụ.

  • Luồng chạy dưới một vai trò, và vai trò đó chỉ đọc được đúng phần dữ liệu được phép. Không có ngoại lệ kiểu “job nền chạy bằng tài khoản quản trị cho tiện”.
  • Cột đã gắn nhãn nhạy cảm giữ nguyên chế độ che khi đi qua luồng.
  • Theo Cloudera công bố, kiểm soát truy cập ở mức chi tiết đi kèm truy vết nguồn gốc tự động, cho hồ sơ xuất xứ dữ liệu từ đầu tới cuối.

Chữ tự động ở dòng cuối là chữ đáng tiền nhất. Sơ đồ luồng dữ liệu vẽ trong tài liệu thiết kế chỉ đúng vào đúng ngày nó được vẽ; sáu tháng sau nó mô tả một hệ thống không còn tồn tại. Truy vết sinh ra từ siêu dữ liệu mà nền tảng ghi lại khi job thật sự chạy thì luôn mô tả hệ thống đang chạy hôm nay. Chi tiết ở Data Catalog & Lineage.

Chia sẻ kết quả: đọc trực tiếp, không sao chép

Phần tiêu đề “Chia sẻ kết quả: đọc trực tiếp, không sao chép”

Luồng chạy xong thì dữ liệu thường phải tới tay người khác — một đơn vị thành viên, một đối tác, một nền tảng phân tích khác. Cách quen thuộc là xuất một bản sao rồi gửi đi, và nó tạo bốn vấn đề cùng lúc: tốn dung lượng, hai nơi lệch số theo thời gian, bản sao nằm ngoài tầm quản trị, và mỗi bản sao là một cửa rò rỉ.

Theo Cloudera công bố, Apache Iceberg cho phép chia sẻ dữ liệu xuyên nền tảng qua REST catalog (REST — Representational State Transfer: kiểu kiến trúc giao tiếp phổ biến trên web), kèm sẵn khả năng tiến hóa cấu trúc bảng và du hành thời gian, dưới sự quản trị của nền tảng.

Cách hiểu gọn: thay vì gửi dữ liệu đi, bạn gửi một đường vào có kiểm soát. Bên nhận đọc chính bảng gốc nên luôn thấy dữ liệu mới nhất; bên chủ dữ liệu giữ nguyên quyền kiểm soát vì bảng vẫn nằm ở chỗ cũ, dưới chính sách cũ. Với các tổ chức tại Việt Nam, điều này chạm đúng một ràng buộc quen thuộc: chia sẻ dữ liệu mà dữ liệu không rời khỏi vành đai đã đăng ký. Xem thêm Open Data Lakehouse & Iceberg.

Nếu đội bạn đã chuẩn hóa phần biến đổi dữ liệu bằng dbt — cách tổ chức logic biến đổi theo mô hình khai báo, có kiểm thử và tài liệu đi kèm — thì hai cách làm bổ trợ nhau rất tự nhiên: Spark lo phần nặng và phần dữ liệu thô, dbt lo tầng mô hình nghiệp vụ phía trên với quy trình rõ ràng. BSD triển khai cả hai và giúp bạn phân vai đúng cho từng tầng.

Một ngân hàng cần bảng giao dịch hợp nhất cho báo cáo sáng hôm sau. Dữ liệu về từ core banking, hệ thống thẻ, kênh số và ví liên kết — mỗi nguồn một khung giờ, một định dạng, và thỉnh thoảng một nguồn về trễ.

Luồng chạy như sau: chờ đủ tín hiệu nguồn (chờ theo phụ thuộc, không theo giờ), nạp thô vào bảng Iceberg theo từng nguồn, Spark chuẩn hóa mã giao dịch và loại bản ghi trùng do gửi lại, đối chiếu tổng số bút toán với hệ thống nguồn, rồi mới dựng bảng phục vụ. Nếu bước đối chiếu lệch, luồng dừng và cảnh báo — bảng phục vụ không bị cập nhật bằng dữ liệu sai. Nếu phát hiện sai sau khi đã ghi, quay lui về ảnh chụp trước đó rồi chạy lại.

Bước đối chiếu là bước hay bị cắt khi tiến độ gấp — và là bước duy nhất đứng giữa một con số sai và một cuộc họp giao ban.

Một khách hàng có thể xuất hiện ở năm hệ thống với năm cách viết tên, hai số điện thoại và ba địa chỉ. Muốn có hồ sơ khách hàng hợp nhất, cần một luồng chuẩn hóa: đưa tên và địa chỉ về dạng thống nhất, chuẩn hóa số điện thoại và mã định danh, so khớp để nhận ra các bản ghi cùng một người, rồi hợp nhất theo thứ tự ưu tiên nguồn.

Việc này chạy lặp lại và luôn cần sửa dần theo thực tế. Ảnh chụp Iceberg giúp mỗi lần đổi luật so khớp đều so được kết quả mới với kết quả cũ trước khi cho áp dụng chính thức.

Chia sẻ: