
Xử Lý Dữ Liệu Thời Gian Thực Với Kafka Streams
Kafka Streams là thư viện xử lý stream thời gian thực của Apache Kafka, cho phép xây dựng ứng dụng phân tích và biến đổi dữ liệu ngay trong ứng dụng Java/Scala mà không cần cluster riêng. Bài viết này hướng dẫn các khái niệm cốt lõi, cách sử dụng DSL, và triển khai production.
Tham khảo: Apache Kafka Streams Documentation
Kafka Streams Là Gì?
Kafka Streams là client library nhúng trong ứng dụng Java/Scala, xử lý dữ liệu từ Kafka topics theo thời gian thực. Không giống Spark Streaming (micro-batch), Kafka Streams xử lý record-by-record với latency millisecond.
- Nhúng (embeddable): Chạy trong cùng JVM với ứng dụng
- Scalable: Tự động phân chia task giữa các instance
- Fault-tolerant: Exactly-once semantics, changelog topics
- Stateful: State stores nội bộ với truy vấn tương tác

Các Khái Niệm Cốt Lõi
Topology
Topology là đồ thị xử lý gồm các stream processors kết nối bằng streams. Mỗi topology đọc từ input topic, xử lý qua các processor, và ghi vào output topic.
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> textLines = builder.stream("input-topic");
KTable<String, Long> wordCounts = textLines
.flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\W+")))
.groupBy((key, value) -> value)
.count();
wordCounts.toStream().to("output-topic", Produced.with(Serdes.String(), Serdes.Long()));
Processor API
Processor API là API thấp cấp cho phép developer viết custom processing logic. Interface Processor<KIn, VIn, KOut, VOut> cung cấp process(record), init(context), và close().
State Stores
State stores lưu trữ dữ liệu trạng thái cho các thao tác stateful (aggregate, join, window). Hỗ trợ persistent key-value store (trên đĩa) và in-memory hashmap. State tự động phục hồi qua changelog topics.
StoreBuilder<KeyValueStore<String, Long>> countStoreBuilder = Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore("counts-store"),
Serdes.String(), Serdes.Long());
Exactly-Once Semantics
Kafka Streams đảm bảo mỗi record được xử lý đúng một lần bằng cách commit atomically offset, state updates, và output records. Exactly-once v2 (từ Kafka 2.6) hiệu quả hơn với throughput cao hơn.
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);

Khái Niệm Thời Gian
Kafka Streams phân biệt ba loại thời gian quan trọng khi xử lý dữ liệu:
- Event time: Thời điểm sự kiện xảy ra tại nguồn, có thể chênh lệch với thời điểm xử lý
- Processing time: Thời điểm record được xử lý bởi ứng dụng Streams
- Ingestion time: Thời điểm record được lưu vào Kafka partition bởi broker
Mặc định, Kafka Streams dùng timestamp của record để xác định event time. Bạn có thể tùy chỉnh qua interface TimestampExtractor nếu dữ liệu nhúng timestamp tùy chỉnh trong payload. Các windowed operations (ví dụ đếm theo giờ) phụ thuộc chính xác vào việc chọn đúng loại thời gian.
Kafka Streams vs Kafka Connect vs Spark Streaming
| Tiêu chí | Kafka Streams | Kafka Connect | Spark Streaming |
|---|---|---|---|
| Mục đích | Xử lý và phân tích real-time | Di chuyển dữ liệu (ETL) | Batch + stream analytics |
| Paradigm | Application library | Connector framework | Distributed framework |
| Latency | Millisecond | N/A (ETL) | Seconds (micro-batch) |
| State | State stores + changelog | Pass-through | RDD/Batch state |
| Fault tolerance | Exactly-once (atomic) | At-least-once | Checkpoint-based |
Các Thao Tác DSL Phổ Biến
Filter — Lọc bản ghi
KStream<String, String> filtered = textLines
.filter((key, value) -> value.contains("error"));
Map — Biến đổi bản ghi
KStream<String, Integer> wordLengths = textLines
.map((key, value) -> new KeyValue<>(value, value.length()));
Join — Kết hợp stream
KStream<String, String> clicks = builder.stream("clicks");
KStream<String, String> impressions = builder.stream("impressions");
KStream<String, String> joined = clicks.join(impressions,
(click, impression) -> click + ":" + impression,
JoinWindows.of(Duration.ofMinutes(5)));
Aggregate — Tổng hợp
KTable<String, Long> wordCounts = textLines
.flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\W+")))
.groupBy((key, value) -> value)
.count();
Interactive Queries
Interactive Queries cho phép ứng dụng bên ngoài query trực tiếp state stores của Kafka Streams mà không cần qua Kafka topics. Điều này hữu ích cho các dashboard real-time, API endpoint phục vụ query điểm (point query) hoặc range query.
- Local state query: Query state store trên instance hiện tại
- Remote state query: Query state store trên instance khác thông qua RPC
- Standalone apps: Ứng dụng riêng chỉ dùng để phục vụ query (không xử lý stream)
// Cấu hình RPC cho Interactive Queries
props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "host:port");
// Query state store từ ứng dụng khác
ReadOnlyKeyValueStore store = streams.store(
StoreQueryParameters.fromNameAndType("counts-store", QueryableStoreTypes.keyValueStore()));
Long count = store.get("kafka");
Triển Khai Production
Cấu Hình Quan Trọng
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "my-streams-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 4);
Monitoring & Security
Sử dụng JMX metrics và StreamsMetrics API để giám sát throughput, latency, state store size. Hỗ trợ SSL/TLS, SASL/PLAIN, SASL/SCRAM cho bảo mật kết nối Kafka.
Các Trường Hợp Sử Dụng Thực Tế
- Real-time Word Count: Đọc log, tách từ, đếm tần suất theo thời gian
- Log Monitoring: Lọc ERROR/FATAL, gửi alert đến topic riêng
- Event-Driven Microservices: Xử lý đơn hàng: xác thực → tính giá → ghi kho → thông báo
- Fraud Detection: Phát hiện giao dịch gian lận dựa trên pattern lịch sử
- Data Enrichment: Kết hợp event stream với KTable reference data
Tham khảo thêm: Developer Guide | Documentation
