引言
流式数据处理是现代数据架构的核心组件。无论是实时指标计算、欺诈检测还是用户行为分析,流处理框架都能提供低延迟、高吞吐的数据处理能力。
流处理核心概念
流处理核心概念:
┌─────────────────────────────────────────┐
│ 事件流(Event Stream) │
│ 无限的事件序列 │
│ 每个事件包含时间戳和数据 │
│ │
│ 时间语义(Time Semantics) │
│ Event Time: 事件实际发生时间 │
│ Processing Time: 事件处理时间 │
│ Ingestion Time: 事件进入系统时间 │
│ │
│ 窗口(Window) │
│ 将无限流切分为有限数据集 │
│ 滚动窗口、滑动窗口、会话窗口 │
│ │
│ 水位线(Watermark) │
│ 衡量事件时间进度的标记 │
│ 处理乱序事件的关键机制 │
│ │
│ 状态(State) │
│ 计算过程中需要持久化的中间结果 │
│ 键控状态、算子状态 │
│ │
│ 精确一次语义(Exactly-Once) │
│ 每个事件精确处理一次 │
│ 通过检查点和事务实现 │
└─────────────────────────────────────────┘
Kafka Streams
基础架构
Kafka Streams架构:
┌─────────────────────────────────────────┐
│ 应用实例1 应用实例2 应用实例3│
│ ┌──────────┐ ┌──────────┐ ┌──────────┐│
│ │Stream任务│ │Stream任务│ │Stream任务││
│ │ 分区0 │ │ 分区1 │ │ 分区2 ││
│ └──────────┘ └──────────┘ └──────────┘│
│ ↓ ↓ ↓ │
│ ┌──────────────────────────────────────┐ │
│ │ 状态存储(RocksDB) │ │
│ └──────────────────────────────────────┘ │
│ ↓ ↓ ↓ │
│ ┌──────────────────────────────────────┐ │
│ │ Kafka Broker │ │
│ │ Topic: input-topic (3 partitions) │ │
│ │ Topic: output-topic (3 partitions) │ │
│ └──────────────────────────────────────┘ │
└─────────────────────────────────────────┘
特点:
- 无需独立集群,作为库嵌入应用
- 自动负载均衡和故障转移
- 支持状态存储和精确一次语义
- 与Kafka生态无缝集成
实时指标计算
package com.example.streams;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import java.time.Duration;
import java.util.Properties;
public class MetricsCalculator {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "metrics-calculator");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG,
Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG,
Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
// 从Kafka Topic读取请求事件
KStream<String, RequestEvent> requests = builder.stream(
"request-events",
Consumed.with(Serdes.String(), new RequestEventSerde())
);
// 1. 实时QPS计算(5秒滚动窗口)
requests
.map((key, event) -> KeyValue.pair(event.getService(), event))
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofSeconds(5)))
.count()
.toStream()
.map((windowedKey, count) -> KeyValue.pair(
windowedKey.key(),
new MetricPoint(
windowedKey.window().start(),
"qps",
count / 5.0
)
))
.to("qps-metrics", Produced.with(Serdes.String(), new MetricPointSerde()));
// 2. 延迟分布计算(1分钟滑动窗口,10秒滑动步长)
requests
.map((key, event) -> KeyValue.pair(event.getService(), event.getLatency()))
.groupByKey()
.windowedBy(SlidingWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(1)))
.aggregate(
LatencyStats::new,
(key, latency, stats) -> stats.add(latency),
(key, stats1, stats2) -> stats1.merge(stats2),
Materialized.with(Serdes.String(), new LatencyStatsSerde())
)
.toStream()
.map((windowedKey, stats) -> KeyValue.pair(
windowedKey.key(),
new LatencyMetric(
windowedKey.window().start(),
stats.getAvg(),
stats.getP50(),
stats.getP95(),
stats.getP99()
)
))
.to("latency-metrics", Produced.with(Serdes.String(), new LatencyMetricSerde()));
// 3. 错误率计算(1分钟滚动窗口)
requests
.map((key, event) -> KeyValue.pair(
event.getService(),
event.getStatus() >= 400 ? 1 : 0
))
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
.aggregate(
() -> new ErrorCount(0, 0),
(key, isError, counts) -> {
counts.total++;
counts.errors += isError;
return counts;
},
Materialized.with(Serdes.String(), new ErrorCountSerde())
)
.toStream()
.map((windowedKey, counts) -> KeyValue.pair(
windowedKey.key(),
new ErrorRate(
windowedKey.window().start(),
(double) counts.errors / counts.total
)
))
.to("error-rate-metrics", Produced.with(Serdes.String(), new ErrorRateSerde()));
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
}
}
class LatencyStats {
private long count = 0;
private double sum = 0;
private final TreeMap<Long, Long> histogram = new TreeMap<>();
public LatencyStats add(long latency) {
count++;
sum += latency;
histogram.merge(latency, 1L, Long::sum);
return this;
}
public LatencyStats merge(LatencyStats other) {
this.count += other.count;
this.sum += other.sum;
other.histogram.forEach((k, v) ->
this.histogram.merge(k, v, Long::sum));
return this;
}
public double getAvg() { return sum / count; }
public long getP50() { return getPercentile(0.50); }
public long getP95() { return getPercentile(0.95); }
public long getP99() { return getPercentile(0.99); }
private long getPercentile(double p) {
long target = (long) (count * p);
long cumulative = 0;
for (var entry : histogram.entrySet()) {
cumulative += entry.getValue();
if (cumulative >= target) {
return entry.getKey();
}
}
return histogram.lastKey();
}
}
会话窗口示例
// 用户会话分析
KStream<String, UserAction> userActions = builder.stream(
"user-actions",
Consumed.with(Serdes.String(), new UserActionSerde())
);
userActions
.groupByKey()
.windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(30)))
.aggregate(
SessionInfo::new,
(key, action, session) -> session.addAction(action),
(key, session1, session2) -> session1.merge(session2),
Materialized.with(Serdes.String(), new SessionInfoSerde())
)
.toStream()
.filter((windowedKey, session) -> session.getDuration() > 60000) // 超过1分钟的会话
.map((windowedKey, session) -> KeyValue.pair(
windowedKey.key(),
new SessionReport(
windowedKey.window().start(),
windowedKey.window().end(),
session.getActionCount(),
session.getDuration(),
session.getPages()
)
))
.to("session-reports", Produced.with(Serdes.String(), new SessionReportSerde()));
Apache Flink
架构对比
Flink vs Kafka Streams:
┌─────────────────────────────────────────┐
│ Apache Flink │
│ ✓ 独立集群,集中管理 │
│ ✓ 强大的窗口和状态管理 │
│ ✓ 支持批处理和流处理 │
│ ✓ 丰富的连接器生态 │
│ ✓ CEP复杂事件处理 │
│ ✗ 部署和运维复杂 │
│ ✗ 资源消耗较大 │
│ │
│ Kafka Streams │
│ ✓ 轻量级,嵌入应用 │
│ ✓ 与Kafka无缝集成 │
│ ✓ 自动扩缩容 │
│ ✓ 运维简单 │
│ ✗ 功能相对较少 │
│ ✗ 仅支持流处理 │
│ │
│ 选择建议: │
│ 简单场景 → Kafka Streams │
│ 复杂计算 → Flink │
│ CEP需求 → Flink │
└─────────────────────────────────────────┘
实时欺诈检测(CEP)
package com.example.flink;
import org.apache.flink.cep.CEP;
import org.apache.flink.cep.PatternSelectFunction;
import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.TumblingEventTimeWindows;
import java.time.Duration;
import java.util.List;
import java.util.Map;
public class FraudDetection {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 启用事件时间
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
// 读取交易事件
DataStream<Transaction> transactions = env
.addSource(new TransactionSource())
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Transaction>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp())
);
// 定义欺诈检测模式
Pattern<Transaction, ?> fraudPattern = Pattern
.<Transaction>begin("first")
.where(new SimpleCondition<Transaction>() {
@Override
public boolean filter(Transaction transaction) {
// 第一笔交易金额大于1000
return transaction.getAmount() > 1000;
}
})
.next("second")
.where(new SimpleCondition<Transaction>() {
@Override
public boolean filter(Transaction transaction) {
// 第二笔交易金额大于2000
return transaction.getAmount() > 2000;
}
})
.next("third")
.where(new SimpleCondition<Transaction>() {
@Override
public boolean filter(Transaction transaction) {
// 第三笔交易金额大于5000
return transaction.getAmount() > 5000;
}
})
.within(Duration.ofMinutes(10)); // 10分钟内完成
// 应用模式检测
DataStream<FraudAlert> fraudAlerts = CEP
.pattern(transactions.keyBy(Transaction::getUserId), fraudPattern)
.select(new PatternSelectFunction<Transaction, FraudAlert>() {
@Override
public FraudAlert select(Map<String, List<Transaction>> pattern) {
Transaction first = pattern.get("first").get(0);
Transaction second = pattern.get("second").get(0);
Transaction third = pattern.get("third").get(0);
return new FraudAlert(
first.getUserId(),
"Suspicious transaction pattern detected",
first.getTimestamp(),
third.getTimestamp(),
first.getAmount() + second.getAmount() + third.getAmount()
);
}
});
// 输出告警
fraudAlerts.print();
fraudAlerts.addSink(new AlertSink());
env.execute("Fraud Detection Job");
}
}
// 另一个模式:异地登录检测
Pattern<LoginEvent, ?>异地登录Pattern = Pattern
.<LoginEvent>begin("login1")
.next("login2")
.where(new SimpleCondition<LoginEvent>() {
@Override
public boolean filter(LoginEvent event) {
// 检查IP地理位置是否变化过大
return true; // 实际应用中需要调用地理位置API
}
})
.within(Duration.ofMinutes(5)); // 5分钟内
状态管理
// 使用键控状态
public class UserSessionCounter extends KeyedProcessFunction<String, UserEvent, SessionCount> {
// 值状态:记录会话数量
private ValueState<Integer> sessionCount;
// 列表状态:记录最近的页面访问
private ListState<String> recentPages;
// 映射状态:记录每个页面的访问次数
private MapState<String, Integer> pageVisits;
// 定时器状态
private ValueState<Long> timerTimestamp;
@Override
public void open(Configuration parameters) {
ValueStateDescriptor<Integer> countDescriptor =
new ValueStateDescriptor<>("session-count", Integer.class);
sessionCount = getRuntimeContext().getState(countDescriptor);
ListStateDescriptor<String> pagesDescriptor =
new ListStateDescriptor<>("recent-pages", String.class);
recentPages = getRuntimeContext().getListState(pagesDescriptor);
MapStateDescriptor<String, Integer> visitsDescriptor =
new MapStateDescriptor<>("page-visits", String.class, Integer.class);
pageVisits = getRuntimeContext().getMapState(visitsDescriptor);
ValueStateDescriptor<Long> timerDescriptor =
new ValueStateDescriptor<>("timer", Long.class);
timerTimestamp = getRuntimeContext().getState(timerDescriptor);
}
@Override
public void processElement(
UserEvent event,
Context ctx,
Collector<SessionCount> out
) throws Exception {
// 更新会话计数
Integer count = sessionCount.value();
if (count == null) {
count = 0;
}
sessionCount.update(count + 1);
// 记录页面访问
recentPages.add(event.getPage());
// 更新页面访问次数
Integer visits = pageVisits.get(event.getPage());
pageVisits.put(event.getPage(), visits == null ? 1 : visits + 1);
// 注册定时器(会话超时)
Long currentTimer = timerTimestamp.value();
long newTimer = ctx.timerService().currentWatermark() + 1800000; // 30分钟
if (currentTimer == null || newTimer > currentTimer) {
if (currentTimer != null) {
ctx.timerService().deleteEventTimeTimer(currentTimer);
}
ctx.timerService().registerEventTimeTimer(newTimer);
timerTimestamp.update(newTimer);
}
}
@Override
public void onTimer(
long timestamp,
OnTimerContext ctx,
Collector<SessionCount> out
) throws Exception {
// 会话超时,输出统计
out.collect(new SessionCount(
ctx.getCurrentKey(),
sessionCount.value(),
StreamSupport.stream(recentPages.get().spliterator(), false)
.collect(Collectors.toList()),
StreamSupport.stream(pageVisits.entries().spliterator(), false)
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue))
));
// 清理状态
sessionCount.clear();
recentPages.clear();
pageVisits.clear();
timerTimestamp.clear();
}
}
Flink SQL 与 Table API
声明式流处理
-- 创建 Kafka 源表
CREATE TABLE user_events (
user_id STRING,
event_type STRING,
amount DECIMAL(10, 2),
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'user-events',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json',
'scan.startup.mode' = 'latest-offset'
);
-- 创建结果输出表(MySQL)
CREATE TABLE event_summary (
event_type STRING PRIMARY KEY NOT ENFORCED,
total_amount DECIMAL(16, 2),
event_count BIGINT,
window_start TIMESTAMP(3)
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/analytics',
'table-name' = 'event_summary',
'username' = 'flink',
'password' = 'password'
);
-- 实时窗口聚合(每 1 分钟计算各事件类型的总金额和次数)
INSERT INTO event_summary
SELECT
event_type,
SUM(amount) AS total_amount,
COUNT(*) AS event_count,
TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start
FROM user_events
GROUP BY
event_type,
TUMBLE(event_time, INTERVAL '1' MINUTE);
Table API 与 DataStream 混合
// 在 DataStream 基础上使用 Table API 做复杂变换
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// DataStream → Table
DataStream<Order> orders = env.addSource(new OrderSource());
tableEnv.createTemporaryView("orders", orders);
// 用 SQL 做聚合
Table result = tableEnv.sqlQuery("""
SELECT user_id, COUNT(*) as cnt, SUM(amount) as total
FROM orders
GROUP BY user_id, TUMBLE(order_time, INTERVAL '5' MINUTES)
HAVING SUM(amount) > 1000
""");
// Table → DataStream,继续流式处理
DataStream<HighValueUser> highValueUsers = tableEnv.toDataStream(result, HighValueUser.class);
highValueUsers.addSink(new AlertSink());
Flink SQL 与传统数仓对比
| 维度 | Flink SQL (实时) | Hive/Spark SQL (离线) |
|---|---|---|
| 数据新鲜度 | 秒级 | 小时到天级 |
| 计算模式 | 增量计算 | 全量扫描 |
| 容错 | Checkpoint 自动恢复 | 失败重跑 |
| 使用成本 | 需理解事件时间和水位线 | SQL 语义与离线一致 |
| 适用 | 实时监控、实时看板 | 报表、T+1 分析 |
生产部署与性能调优
Kubernetes 部署模式
# flink-deployment.yaml
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: fraud-detection-pipeline
spec:
image: flink:1.18-scala_2.12
flinkVersion: v1.18
serviceAccount: flink-sa
jobManager:
resource:
memory: "2Gi"
cpu: 1
taskManager:
resource:
memory: "4Gi"
cpu: 2
replicas: 3
job:
jarURI: local:///opt/flink/usrlib/fraud-detection.jar
parallelism: 6
upgradeMode: savepoint
state: running
性能调优检查清单
| 指标 | 问题 | 调优方案 |
|---|---|---|
| 延迟高 | Checkpoint 间隔太短 | 增大 checkpointing.interval 至 30s-60s |
| 背压 | 下游消费慢 | 增大并行度 / 优化算子 / 增加资源 |
| OOM | 状态过大 | 启用增量 Checkpoint + 状态 TTL |
| GC 频繁 | 堆内存不足 | 增大 TaskManager 内存 / 调整 G1GC |
| 序列化慢 | 使用 Java 序列化 | 切换为 Avro / Protobuf / Kryo |
// 优化:启用增量 Checkpoint + 本地恢复
env.enableCheckpointing(60000);
env.getCheckpointConfig().setCheckpointStorage("hdfs://namenode:8020/flink/checkpoints");
env.getCheckpointConfig().enableIncrementalCheckpoints(true);
env.getCheckpointConfig().setPreferCheckpointForRecovery(true);
// 优化:状态 TTL 自动清理过期数据
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.cleanupIncrementally(10, true)
.build();
总结
流处理框架选型
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| 简单ETL | Kafka Streams | 轻量,易部署 |
| 实时指标 | Kafka Streams | 与Kafka集成好 |
| 复杂CEP | Flink CEP | 强大的模式匹配 |
| 批流一体 | Flink | 统一计算引擎 |
| 机器学习 | Flink ML | 内置算法库 |
| SQL分析 | Flink SQL | 声明式查询 |
关键原则
- 选择合适的语义:At-Least-Once vs Exactly-Once
- 合理设置窗口:平衡延迟和准确性
- 管理状态大小:定期清理过期状态
- 监控背压:防止慢消费者拖垮系统
- 测试乱序处理:验证水位线机制
- 优化序列化:选择高效的序列化格式
- 资源规划:根据数据量分配内存和CPU
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。