一、Kafka Streams概述
Kafka Streams是Apache Kafka的流处理库,用于构建实时流处理应用。Kafka Streams提供了丰富的API,支持状态管理、窗口操作、聚合计算等功能,是构建实时数据管道的理想选择。
二、Kafka Streams架构
2.1 架构设计
graph TD
A[Kafka Topic] --> B[Kafka Streams应用]
B --> C[Source Processor]
C --> D[Stream Processor]
D --> E[Sink Processor]
E --> F[Kafka Topic]
B --> G[状态存储]
D --> G
2.2 核心概念
| 概念 | 说明 | 作用 |
|---|---|---|
| Stream | 无界数据流 | 处理实时数据 |
| Table | 有界数据集 | 状态管理 |
| Processor | 数据处理单元 | 转换数据 |
| State Store | 状态存储 | 存储中间结果 |
| Window | 时间窗口 | 窗口聚合 |
三、DSL API编程
3.1 创建Streams应用
StreamsBuilder builder = new StreamsBuilder();
// 从Kafka Topic读取数据
KStream<String, String> stream = builder.stream("input-topic");
// 处理数据
stream
.mapValues(value -> value.toUpperCase())
.filter((key, value) -> value.contains("IMPORTANT"))
.to("output-topic");
// 创建拓扑
Topology topology = builder.build();
// 配置Streams
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "my-streams-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
// 启动应用
KafkaStreams streams = new KafkaStreams(topology, props);
streams.start();
3.2 数据转换操作
// Map操作
stream.map((key, value) -> KeyValue.pair(key, value.toUpperCase()));
stream.mapValues(value -> value.length());
// Filter操作
stream.filter((key, value) -> value.length() > 10);
stream.filterNot((key, value) -> value.isEmpty());
// SelectKey操作
stream.selectKey((key, value) -> extractUserId(value));
// FlatMap操作
stream.flatMapValues(value -> Arrays.asList(value.split(",")));
3.3 聚合操作
// GroupBy操作
KGroupedStream<String, String> grouped = stream.groupByKey();
// Count聚合
KTable<String, Long> countTable = grouped.count();
// Sum聚合
KTable<String, Long> sumTable = grouped.aggregate(
() -> 0L,
(key, value, aggregate) -> aggregate + value.length(),
Materialized.as("sum-store")
);
// Reduce聚合
KTable<String, String> reduceTable = grouped.reduce(
(value1, value2) -> value1 + "," + value2
);
四、窗口操作
4.1 时间窗口类型
graph LR
A[时间线] --> B[窗口1: 0-10s]
A --> C[窗口2: 10-20s]
A --> D[窗口3: 20-30s]
E[事件1: t=5s] --> B
F[事件2: t=15s] --> C
G[事件3: t=25s] --> D
4.2 滚动窗口
// 滚动窗口:每10秒一个窗口
TimeWindowedKStream<String, String> windowed = stream
.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofSeconds(10)));
KTable<Windowed<String>, Long> countTable = windowed.count();
4.3 滑动窗口
// 滑动窗口:窗口大小30秒,滑动间隔5秒
TimeWindowedKStream<String, String> windowed = stream
.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofSeconds(30)).advanceBy(Duration.ofSeconds(5)));
KTable<Windowed<String>, Long> countTable = windowed.count();
4.4 会话窗口
// 会话窗口:空闲时间5分钟后结束会话
SessionWindowedKStream<String, String> sessionWindowed = stream
.groupByKey()
.windowedBy(SessionWindows.with(Duration.ofMinutes(5)));
KTable<Windowed<String>, Long> countTable = sessionWindowed.count();
五、状态管理
5.1 状态存储类型
| 存储类型 | 说明 | 适用场景 |
|---|---|---|
| KeyValueStore | 键值存储 | 聚合状态 |
| WindowStore | 窗口存储 | 窗口聚合 |
| SessionStore | 会话存储 | 会话分析 |
5.2 自定义状态存储
// 使用自定义状态存储
KTable<String, Long> countTable = stream
.groupByKey()
.count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("custom-store")
.withKeySerde(Serdes.String())
.withValueSerde(Serdes.Long()));
5.3 状态恢复
Kafka Streams通过变更日志(Changelog)自动恢复状态:
graph TD
A[状态存储] --> B[变更日志Topic]
B --> C[持久化到Kafka]
D[应用重启] --> E[从变更日志恢复]
E --> A
六、Stream-Table Join
6.1 Stream-Stream Join
KStream<String, Order> orderStream = builder.stream("orders");
KStream<String, Payment> paymentStream = builder.stream("payments");
// Stream-Stream Join,时间窗口5分钟
KStream<String, OrderPayment> joinedStream = orderStream
.join(paymentStream,
(order, payment) -> new OrderPayment(order, payment),
JoinWindows.of(Duration.ofMinutes(5)));
6.2 Stream-Table Join
KStream<String, Order> orderStream = builder.stream("orders");
KTable<String, Customer> customerTable = builder.table("customers");
// Stream-Table Join
KStream<String, OrderWithCustomer> joinedStream = orderStream
.leftJoin(customerTable,
(order, customer) -> new OrderWithCustomer(order, customer));
6.3 Table-Table Join
KTable<String, Customer> customerTable = builder.table("customers");
KTable<String, Address> addressTable = builder.table("addresses");
// Table-Table Join
KTable<String, CustomerWithAddress> joinedTable = customerTable
.join(addressTable,
(customer, address) -> new CustomerWithAddress(customer, address));
七、Kafka Streams与其他流处理框架对比
| 框架 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| Kafka Streams | 轻量、集成Kafka | 功能有限 | 中小型流处理 |
| Flink | 功能强大 | 复杂 | 大型流处理 |
| Spark Streaming | 批流一体 | 微批延迟 | 批流混合 |
八、Kafka Streams性能优化
8.1 分区优化
// 合理设置Topic分区数
// 分区数应等于或大于应用实例数
8.2 缓存优化
// 配置缓存大小
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "1000");
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, "1024000"); // 1MB
8.3 状态存储优化
// 使用RocksDB作为状态存储
props.put(StreamsConfig.DEFAULT_STATE_STORE_CONFIG,
RocksDBConfig.newBuilder().setNumThreads(4).build());
8.4 并行度配置
// 设置并行度
// 每个分区由一个线程处理
// 增加分区数可以提高并行度
九、Kafka Streams部署
9.1 独立部署
// 打包成Jar包运行
java -jar my-streams-app.jar
9.2 Docker部署
# Dockerfile
FROM openjdk:11
COPY my-streams-app.jar /app/
CMD ["java", "-jar", "/app/my-streams-app.jar"]
9.3 Kubernetes部署
# deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: kafka-streams-app
spec:
replicas: 3
selector:
matchLabels:
app: kafka-streams
template:
metadata:
labels:
app: kafka-streams
spec:
containers:
- name: kafka-streams
image: my-streams-app:latest
十、总结
Kafka Streams是一款轻量级的流处理框架,与Kafka深度集成,适合构建实时数据管道。掌握Kafka Streams的核心概念和API,能够快速构建高效的实时流处理应用,满足数据密集型系统的实时处理需求。