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

Cloudera Streaming: Apache Kafka, Apache Flink & Edge Management

Có loại dữ liệu mà giá trị của nó rụng theo từng giây. Một giao dịch đáng ngờ phát hiện sau hai tiếng chỉ còn là số liệu để báo cáo; phát hiện trong 200 mili giây thì chặn được. Một máy nén rung bất thường báo vào sáng hôm sau là biên bản sự cố; báo ngay lúc đó là một lần bảo trì có kế hoạch. Cloudera Streaming dành cho đúng nhóm bài toán này: xử lý dữ liệu khi nó vẫn đang xảy ra.

Nền của nó là hai dự án mã nguồn mở với hai vai trò hoàn toàn khác nhauApache Kafka lo việc nhận, đệm và phân phối sự kiện; Apache Flink lo việc tính toán trên dòng sự kiện đó. Kèm theo là Edge Management để lấy dữ liệu từ thiết bị ở biên. Lẫn vai trò của Kafka và Flink là nguồn gốc của phần lớn thiết kế sai, nên bài này tách bạch trước.

flowchart LR
  subgraph EDGE["Biên"]
    E1["Cảm biến nhà máy"]
    E2["Thiết bị · trạm viễn thông"]
  end
  A["Ứng dụng nghiệp vụ<br/>giao dịch · log"]
  subgraph CSP["Cloudera Streaming"]
    K["Apache Kafka<br/>nhận · đệm · phân phối sự kiện"]
    F["Apache Flink<br/>tính toán có trạng thái<br/>cửa sổ thời gian"]
    K --> F
    F -.->|"kết quả quay lại dòng"| K
  end
  subgraph OUT["Đích"]
    O1["Cảnh báo tức thì<br/>chặn giao dịch"]
    O2["Operational Database<br/>tra cứu thời gian thực"]
    O3["Lakehouse<br/>bảng Iceberg"]
  end
  E1 --> K
  E2 --> K
  A --> K
  F --> O1
  F --> O2
  F --> O3

Cách nhớ gọn nhất: Kafka là băng chuyền có kho đệm. Flink là trạm xử lý đặt trên băng chuyền đó. Băng chuyền không biết món hàng có gì bất thường; trạm xử lý không có việc gì làm nếu thiếu băng chuyền.

Apache KafkaApache Flink
Trả lời câu hỏiSự kiện nằm ở đâu, ai đọc được, giữ bao lâuSự kiện này có ý nghĩa gì khi đặt cạnh những sự kiện trước nó
Vai tròHàng đợi sự kiện phân tán: nhận, lưu bền, phân phốiMáy xử lý dòng: tính toán liên tục trên dòng sự kiện
Giữ cái gìGiữ dữ liệu, theo thời hạn được cấu hìnhGiữ trạng thái tính toán: đếm, cộng dồn, phiên hoạt động
Khái niệm lõiTopic, partition, consumer group, retentionCửa sổ thời gian, trạng thái, watermark, checkpoint
Thiếu nó thì saoNguồn phải nối thẳng vào từng ứng dụng, một đầu chết là gãy cả dâyCó dòng sự kiện chạy qua nhưng không ai rút ra kết luận

Apache Kafka — xương sống của dòng sự kiện

Phần tiêu đề “Apache Kafka — xương sống của dòng sự kiện”

Kafka nhận sự kiện từ nhiều nguồn, ghi bền xuống đĩa rồi cho nhiều hệ thống cùng đọc lại. Ba tính chất khiến nó thành lớp nền:

  • Tách rời người gửi và người nhận. Hệ thống lõi chỉ cần bắn sự kiện vào một topic. Sau đó ai đọc — đội chống gian lận, đội báo cáo, đội chăm sóc khách hàng — không còn là việc của nó. Thêm một hệ thống tiêu thụ mới không phải sửa nguồn.
  • Mỗi bên đọc theo nhịp riêng. Cùng một dòng sự kiện, đội này đọc tức thì để cảnh báo, đội kia đọc chậm hơn để nạp vào kho. Không ai làm chậm ai.
  • Giữ lại và phát lại. Sự kiện nằm trong topic theo thời hạn cấu hình. Phát hiện logic sai hôm qua, bạn cho chạy lại từ đầu dòng thay vì đi xin lại dữ liệu từ hệ thống nguồn.

Kafka cũng chính là tấm đệm chịu tải đỉnh: lúc cao điểm sự kiện dồn vào topic, phía sau xử lý theo sức của mình, không ai vỡ.

Nếu Kafka trả lời “sự kiện gì vừa tới”, Flink trả lời “chuyện gì đang diễn ra”. Khác biệt nằm ở hai chữ có trạng thái: Flink nhớ những gì đã đi qua, nên nó nhìn được chuỗi chứ không chỉ từng điểm rời rạc.

Trạng thái là phần Flink giữ lại giữa các sự kiện — số lần rút tiền của một khách trong hôm nay, tổng hạn mức đã dùng, nhiệt độ trung bình của một máy trong mười phút qua. Cửa sổ thời gian là cách bạn cắt dòng vô tận thành từng lát để tính:

  • Cửa sổ trượt — “năm phút gần nhất”, cập nhật liên tục. Hợp với phát hiện bất thường.
  • Cửa sổ cố định — “từng khung mười lăm phút”. Hợp với thống kê, đối soát.
  • Cửa sổ theo phiên — gom các sự kiện của cùng một lần hoạt động, tự đóng khi khách ngừng thao tác.

Thời điểm sự kiện, không phải thời điểm nhận được

Phần tiêu đề “Thời điểm sự kiện, không phải thời điểm nhận được”

Dữ liệu về trễ là chuyện thường: sóng yếu, thiết bị mất kết nối rồi gửi bù. Flink tính theo thời điểm sự kiện thực sự xảy ra, dùng cơ chế watermark để biết khi nào một cửa sổ có thể đóng và vẫn đón được bản ghi tới muộn. Nhờ vậy con số ra đúng với những gì đã diễn ra ngoài đời, không phải đúng với thứ tự gói tin về tới máy chủ.

Cùng với đó, checkpoint cho phép Flink khôi phục cả trạng thái khi có sự cố — số liệu không bị đếm hụt cũng không bị đếm hai lần.

Viết bằng SQL, không nhất thiết phải viết mã

Phần tiêu đề “Viết bằng SQL, không nhất thiết phải viết mã”

Cloudera đóng gói kèm SQL Stream Builder — bạn khai báo logic dòng bằng SQL (Structured Query Language) quen thuộc thay vì viết ứng dụng Flink từ đầu. Người phân tích nghiệp vụ đọc được luật đang chạy; luật sửa được trong ngày thay vì chờ một chu kỳ phát hành phần mềm.

Edge Management — lấy dữ liệu từ nơi nó sinh ra

Phần tiêu đề “Edge Management — lấy dữ liệu từ nơi nó sinh ra”

Nhiều dữ liệu có giá trị nhất lại sinh ra ở chỗ xa nhất: máy trong xưởng, trạm phát sóng, thiết bị IoT (Internet of Things) ngoài hiện trường. Cloudera Edge Management đặt một tác tử rất nhẹ ngay trên thiết bị, và một trung tâm để điều khiển hàng nghìn tác tử đó từ xa — đẩy luồng mới xuống, cập nhật luật lọc, theo dõi tình trạng.

Điểm quan trọng: tác tử lọc và gộp ngay tại chỗ. Cảm biến sinh dữ liệu mỗi giây nhưng chỉ giá trị vượt ngưỡng mới đáng gửi về. Xử lý ở biên trước khi truyền giúp đường truyền và chi phí lưu trữ giảm đi rõ rệt — với hiện trường ở Việt Nam, nơi đường truyền không phải chỗ nào cũng khỏe, đây thường là điều kiện để dự án chạy được.

Bài toánKafka làm gìFlink làm gì
Phát hiện gian lận giao dịchNhận dòng giao dịch thẻ và chuyển khoản từ nhiều kênh, cho nhiều đội cùng đọcGiữ trạng thái theo từng khách trong cửa sổ trượt vài phút, bắt chuỗi bất thường: nhiều lần thử liên tiếp, hai giao dịch ở hai nơi cách xa nhau trong thời gian không hợp lý
Cảnh báo hạn mứcĐệm sự kiện phát sinh, chịu tải đỉnh giờ cao điểmCộng dồn theo khách theo ngày, vượt ngưỡng thì phát cảnh báo ngay, không chờ chốt sổ cuối ngày
Cảm biến nhà máyNhận dữ liệu do tác tử ở biên gửi về sau khi đã lọcTính trung bình trượt nhiệt độ và độ rung, phát hiện xu hướng lệch trước khi thành hỏng hóc
Dữ liệu viễn thôngGánh khối lượng rất lớn bản ghi CDR (Call Detail Record) từ toàn mạngTổng hợp chất lượng theo trạm theo cửa sổ thời gian, phát hiện vùng suy giảm

Trong cả bốn, sự kiện thường đi tiếp vào lakehouse để phân tích dài hạn — cùng một dòng, vừa dùng ngay vừa để dành. Và vì Streaming nằm trong nền tảng Cloudera, quyền truy cập với dòng dữ liệu này áp theo SDX (Shared Data Experience) như mọi dịch vụ khác: chính sách đặt một lần, không có ngoại lệ cho dữ liệu thời gian thực.

Khi nào dùng Data Flow, khi nào dùng Streaming?

Phần tiêu đề “Khi nào dùng Data Flow, khi nào dùng Streaming?”

Hai lớp này bổ trợ nhau, không thay nhau. Chọn theo bản chất công việc:

Việc cần làmNghiêng về
Nối nhiều nguồn và đích khác kiểu, gồm cả giao thức cũData Flow
Xử lý từng bản ghi độc lập: lọc, đổi định dạng, định tuyến, làm giàuData Flow
Cần xuất xứ từng bản ghi và phát lại đúng bản ghi đóData Flow
Nhiều hệ thống cùng cần đọc một dòng sự kiện, mỗi bên một nhịpKafka
Cần đệm chịu tải đỉnh, hoặc phát lại cả dòng sau khi sửa lỗiKafka
Cộng dồn, ghép dòng, bắt chuỗi hành vi theo cửa sổ thời gianFlink
Độ trễ tính bằng giây hoặc dưới giây, chạy liên tụcFlink

Trong hệ thống thật, câu trả lời gần như luôn là cả hai, xếp thành một chuỗi: Data Flow thu gom từ những nguồn khó và đẩy vào Kafka; Kafka đệm và phân phối; Flink tính toán ra kết luận; rồi kết quả quay lại Kafka hoặc theo Data Flow đi tới các đích cuối. Mỗi lớp làm đúng phần mình mạnh nhất.

  • Dùng Kafka như kho lưu trữ lâu dài. Kafka giỏi giữ dữ liệu nóng cho việc phát lại; dữ liệu để dành phân tích lâu dài thuộc về lakehouse.
  • Giao cho Flink việc di chuyển tệp. Kéo tệp từ máy chủ cũ về là việc của Data Flow — dùng máy xử lý dòng để làm việc đó thì vừa nặng vừa khó nhìn.
  • Đặt thời hạn giữ quá ngắn. Tới lúc phát hiện logic sai mới thấy dòng cũ đã hết hạn — mất luôn khả năng chạy lại.
  • Bỏ qua quản lý schema. Nguồn đổi một trường, mọi bên tiêu thụ gãy cùng lúc, và không ai biết đổi từ lúc nào.
  • Tính theo thời điểm nhận được thay vì thời điểm sự kiện xảy ra. Dữ liệu về trễ sẽ rơi nhầm cửa sổ, con số lệch mà nhìn vẫn thấy hợp lý.
  • Gửi thô mọi thứ từ biên về trung tâm. Phần lớn lưu lượng đó không ai đọc, nhưng vẫn phải trả tiền đường truyền và lưu trữ.
Chia sẻ: