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

Cloudera Data Flow: chuyển dữ liệu mọi nguồn với Apache NiFi

Phần khó nhất của một dự án dữ liệu hiếm khi là mô hình phân tích. Nó là việc lấy được dữ liệu ra khỏi nơi nó đang nằm: máy chủ ở chi nhánh, hệ thống lõi chạy đã mười mấy năm, tệp đẩy qua FTP (File Transfer Protocol) lúc nửa đêm, thiết bị đặt ngoài xưởng. Cloudera Data Flow sinh ra cho đúng phần việc đó — thu thập và chuyển dữ liệu từ mọi nguồn tới mọi đích, an toàn và tiết kiệm.

Nền của nó là Apache NiFi, dự án mã nguồn mở chuyên về tự động hóa dòng dữ liệu. Thay vì mỗi đường dữ liệu một script riêng, bạn vẽ luồng trên giao diện: kéo các khối xử lý ra, nối lại, đặt điều kiện định tuyến. Luồng vừa chạy được vừa đọc được — người vận hành nhìn vào là biết dữ liệu đi đường nào, đang tắc ở đâu.

flowchart LR
  subgraph SRC["Nguồn"]
    S1["Hệ thống lõi<br/>ngân hàng"]
    S2["Tệp trên FTP<br/>máy chủ chi nhánh"]
    S3["Thiết bị · cảm biến<br/>hàng đợi tin nhắn"]
  end
  subgraph CDF["Cloudera Data Flow — Apache NiFi"]
    direction LR
    P1["Thu thập"] --> P2["Định tuyến<br/>lọc · làm giàu"]
    P2 --> P3["Biến đổi<br/>chuẩn hóa định dạng"]
    P3 --> P4["Phân phối"]
  end
  subgraph DST["Đích"]
    D1["Lakehouse<br/>bảng Iceberg"]
    D2["Kafka"]
    D3["Kho object · dịch vụ khác"]
  end
  S1 --> P1
  S2 --> P1
  S3 --> P1
  P4 --> D1
  P4 --> D2
  P4 --> D3

Luồng kéo-thả — dữ liệu đi đường nào, nhìn là thấy

Phần tiêu đề “Luồng kéo-thả — dữ liệu đi đường nào, nhìn là thấy”

Bốn khái niệm sau đủ để đọc hiểu bất kỳ luồng NiFi nào:

Khái niệmLà gìVai trò trong luồng
FlowFileMột đơn vị dữ liệu đang di chuyển: phần nội dung tập thuộc tính đi kèmThứ chảy trong luồng — mỗi bản ghi, mỗi tệp là một FlowFile
ProcessorMột khối làm một việc: đọc, lọc, đổi định dạng, gọi API, ghi ra đíchNút xử lý bạn kéo ra màn hình
ConnectionHàng đợi nối hai processor, có ngưỡng và thứ tự ưu tiênNơi dữ liệu chờ, cũng là nơi phát hiện tắc
Process GroupNhóm nhiều processor thành một khối có đầu vào và đầu raChia luồng lớn thành phần việc gọn, giao cho từng đội

Cloudera cung cấp sẵn hàng trăm processor cho các nguồn và đích phổ biến (theo Cloudera công bố), kèm bộ luồng dựng sẵn cho những đường dữ liệu hay gặp. Luồng được lưu phiên bản vào kho luồng — sửa gì, ai sửa, quay lại bản cũ đều rõ ràng. Đây là khác biệt giữa một đường dữ liệu có thể bàn giao và một script chỉ tác giả hiểu.

Đây là thế mạnh riêng của nền NiFi. Với mỗi FlowFile, hệ thống ghi lại toàn bộ chuỗi sự kiện: nó vào luồng lúc nào, từ đâu, đi qua processor nào, bị đổi thành gì, cuối cùng rơi vào đích nào — kèm dấu thời gian. Đó là data provenance (xuất xứ dữ liệu).

Ba việc mà năng lực này giải quyết:

  • Trả lời câu hỏi “số này ở đâu ra”: lần ngược từ bản ghi trong lakehouse về đúng tệp gốc, đúng lần chạy.
  • Điều tra sự cố: một bản ghi sai định dạng đi tới đâu thì hỏng, thay vì đoán mò cả luồng.
  • Phát lại (replay): chạy lại đúng bản ghi đó qua luồng sau khi vá lỗi, không cần kéo lại toàn bộ dữ liệu nguồn.

Với ngân hàng và khu vực công — nơi câu hỏi kiểm toán không phải “hệ thống có chạy không” mà “chứng minh con số này đúng nguồn” — hồ sơ xuất xứ này là bằng chứng có sẵn, không phải thứ dựng lại sau.

Nguồn dữ liệu không bao giờ hỏi đích còn kham nổi hay không. Đích chậm lại, hoặc bảo trì vài giờ, mà nguồn vẫn bơm đều — đó là lúc luồng thường vỡ.

NiFi xử lý bằng back-pressure (xử lý ngược áp lực): mỗi hàng đợi có ngưỡng theo số bản ghi và theo dung lượng. Chạm ngưỡng, processor phía trước tự ngừng nhận thêm thay vì cố đẩy. Áp lực dội ngược lên đầu luồng một cách có trật tự; khi đích khỏe lại, luồng tự chảy tiếp.

Đi kèm là vài cơ chế nhỏ nhưng quan trọng trong vận hành thật:

  • Ưu tiên hàng đợi: dữ liệu quan trọng đi trước khi ùn.
  • Hết hạn dữ liệu: bản ghi quá cũ tự rụng, không đọng mãi.
  • Bảo đảm giao nhận: hệ thống ghi nhật ký trước khi xử lý, mất điện giữa chừng thì khôi phục đúng chỗ đang dở.

Kết quả thực tế: sự cố ở đích thành chậm trễ có kiểm soát, không thành mất dữ liệu.

Dữ liệu lúc di chuyển là lúc dễ tổn thương nhất. Cloudera Data Flow đóng ba lớp:

  • Mã hóa đường truyền bằng TLS (Transport Layer Security) cho cả kết nối vào nguồn, ra đích và giữa các nút trong cụm.
  • Phân quyền tới từng luồng theo mô hình RBAC (Role-Based Access Control): đội này chỉ thấy và sửa luồng của mình.
  • Che dữ liệu nhạy cảm ngay trong luồng: che số thẻ, số căn cước trước khi bản ghi rời khỏi vành đai an toàn.

Vì Data Flow nằm trong nền tảng Cloudera, các chính sách này dùng chung SDX (Shared Data Experience) với phần còn lại — đặt một lần, áp cho cả kho, cho AI, cho luồng.

Hàng chục điểm phát sinh dữ liệu, đường truyền không đều, có nơi rớt mạng vài tiếng là chuyện thường. Luồng đặt ở mỗi vùng sẽ gom, nén, gắn nhãn nguồn rồi đẩy về trung tâm; chỗ nào rớt mạng thì dữ liệu nằm chờ trong hàng đợi và tự đi tiếp khi có sóng lại — không ai phải chạy tay bù dữ liệu ngày hôm sau.

Đẩy dữ liệu từ hệ thống lõi vào lakehouse

Phần tiêu đề “Đẩy dữ liệu từ hệ thống lõi vào lakehouse”

Hệ thống lõi ngân hàng không chịu nổi việc bị quét toàn bảng giờ hành chính. Cách làm gọn là CDC (Change Data Capture) — chỉ bắt phần thay đổi — rồi để luồng chuẩn hóa và ghi vào bảng Iceberg trong lakehouse. Tải lên hệ thống lõi gần như không đổi, còn phía phân tích luôn có dữ liệu mới.

Thu dữ liệu từ hệ thống cũ không có API hiện đại

Phần tiêu đề “Thu dữ liệu từ hệ thống cũ không có API hiện đại”

Đây là phần việc thực tế Việt Nam gặp mỗi ngày: tệp CSV trên thư mục chia sẻ, bảng dữ liệu trong cơ sở dữ liệu đời cũ, hàng đợi tin nhắn của hệ thống mua từ mười năm trước, tệp XML sinh theo lô ban đêm. Data Flow nói được các giao thức cũ đó và biến chúng thành dòng dữ liệu chuẩn — bạn hiện đại hóa lớp dữ liệu mà không phải viết lại hệ thống nghiệp vụ đang chạy ổn.

Ba chỗ, đều là tiền thật:

  1. Lọc và giảm ngay tại điểm thu. Bỏ trường thừa, gộp bản ghi vụn ngay ở đầu nguồn thì băng thông và dung lượng lưu trữ ở đích giảm theo — thay vì chở hết về rồi mới lọc.
  2. Không trả tiền cho việc chạy lại. Có back-pressure và bảo đảm giao nhận, số lần phải kéo lại dữ liệu vì một sự cố nhỏ giảm hẳn.
  3. Không nuôi một đội chỉ để trông script. Luồng có phiên bản, có cảnh báo, có xuất xứ — thời gian đội dữ liệu quay lại cho việc tạo giá trị.
  • Dựng luồng cho tới lúc “chạy được” rồi dừng. Ngưỡng back-pressure để mặc định thì hôm đích chậm mới biết là chưa đủ.
  • Chở hết về rồi lọc sau. Phần lớn chi phí đường truyền và lưu trữ đến từ dữ liệu không ai dùng.
  • Không lưu phiên bản luồng. Luồng sửa nóng trên môi trường chạy thật sẽ thành hộp đen sau vài tháng.
  • Chỉ nhớ tới xuất xứ khi kiểm toán hỏi. Bật sẵn thì nó là hồ sơ; bật sau thì quá khứ đã mất.
  • Giao cho Data Flow việc tính toán có trạng thái. Cộng dồn theo cửa sổ thời gian, ghép nhiều dòng sự kiện — đó là việc của lớp Streaming, để đúng chỗ thì cả hai đều nhẹ.
Chia sẻ: