DataStream API 是 Apache Flink 提供的底层流处理 API。它允许细致控制状态、时间和自定义处理逻辑,适合构建复杂的事件驱动应用。
要构建的应用
本教程通过信用卡交易示例,构建一个识别特定交易模式的教学程序。数据流如下:
Transactions (source) → Flink (KeyedProcessFunction) → Alerts (sink)
你将学习设置执行环境、创建数据源、用 keyBy 分区以支持并行处理、用 KeyedProcessFunction 实现业务逻辑,以及用托管的 ValueState 保存每个账户的状态。
前置知识与求助
最好了解 Java 或 Python;熟悉其他语言也可以跟随步骤学习。遇到问题时可查阅 社区支持资源,或通过 用户邮件列表求助。
准备开发环境
原教程列出的环境如下:
- Java 路径:Java 11、17 或 21,以及 Maven。
- Python 路径:Java 11、17 或 21,以及 Python 3.9、3.10、3.11 或 3.12。
版本补充:同版本的 Java 兼容性说明推荐 Java 17;Java 21 支持标为实验性。配置环境时应同时检查该说明及所用连接器的限制,而不是把教程的简表理解为所有组合均无条件支持。
Java 项目
Flink 提供 Maven archetype,可快速创建骨架项目,把注意力放在业务逻辑上。依赖包括流应用的核心库 flink-streaming-java,以及提供示例数据生成器等类的 flink-walkthrough-common。
$ mvn archetype:generate \
-DarchetypeGroupId=org.apache.flink \
-DarchetypeArtifactId=flink-walkthrough-datastream-java \
-DarchetypeVersion=2.3.0 \
-DgroupId=frauddetection \
-DartifactId=frauddetection \
-Dversion=0.1 \
-Dpackage=frauddetection \
-DinteractiveMode=false
可以按需修改 groupId、artifactId 和 package。以上参数创建名为 frauddetection 的目录,其中包含完成本教程所需的依赖。Maven archetype 仅供已发布版本使用;这里的 2.3.0 来自所读 2.3 文档的实际渲染值。
把项目导入编辑器后,可找到 FraudDetectionJob.java,并在 IDE 中运行。若遇到 java.lang.NoClassDefFoundError,可能是 classpath 缺少必要的 Flink 依赖。在 IntelliJ IDEA 中依次打开 Run → Edit Configurations → Modify options,并选择 “include dependencies with ‘Provided’ scope”。
Python 项目
Python DataStream API 需要 PyFlink。它可从 PyPI安装:
$ python -m pip install apache-flink
建议使用 虚拟环境隔离依赖。安装后创建 fraud_detection.py。原命令没有固定版本;要复现本教程,应核对安装的 PyFlink 版本与 2.3 文档对应关系。这里保留原命令,未进行安装。
完整程序概览
先看 Java 与 Python 的完整示例,再分解各部分:
Java
public class FraudDetectionJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<Transaction> transactions = env
.fromSource(
TransactionSource.unbounded(),
WatermarkStrategy.noWatermarks(),
"transactions");
DataStream<Alert> alerts = transactions
.keyBy(Transaction::getAccountId)
.process(new FraudDetector())
.name("fraud-detector");
alerts
.addSink(new AlertSink())
.name("send-alerts");
env.execute("Fraud Detection");
}
}
public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {
private static final long serialVersionUID = 1L;
private static final double SMALL_AMOUNT = 1.00;
private static final double LARGE_AMOUNT = 500.00;
private transient ValueState<Boolean> flagState;
@Override
public void open(OpenContext openContext) {
ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
"flag",
Types.BOOLEAN);
flagState = getRuntimeContext().getState(flagDescriptor);
}
@Override
public void processElement(
Transaction transaction,
Context context,
Collector<Alert> collector) throws Exception {
Boolean lastTransactionWasSmall = flagState.value();
if (lastTransactionWasSmall != null) {
if (transaction.getAmount() > LARGE_AMOUNT) {
Alert alert = new Alert();
alert.setId(transaction.getAccountId());
collector.collect(alert);
}
flagState.clear();
}
if (transaction.getAmount() < SMALL_AMOUNT) {
flagState.update(true);
}
}
}
Python
from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.functions import KeyedProcessFunction, RuntimeContext
from pyflink.datastream.state import ValueStateDescriptor
class FraudDetector(KeyedProcessFunction):
SMALL_AMOUNT = 1.00
LARGE_AMOUNT = 500.00
def __init__(self):
self.flag_state = None
def open(self, runtime_context: RuntimeContext):
descriptor = ValueStateDescriptor("flag", Types.BOOLEAN())
self.flag_state = runtime_context.get_state(descriptor)
def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
# transaction is a tuple: (account_id, timestamp, amount)
account_id = transaction[0]
amount = transaction[2]
last_transaction_was_small = self.flag_state.value()
if last_transaction_was_small is not None:
if amount > self.LARGE_AMOUNT:
yield f"Alert{{id={account_id}}}"
self.flag_state.clear()
if amount < self.SMALL_AMOUNT:
self.flag_state.update(True)
def fraud_detection():
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)
# Sample transaction data: (account_id, timestamp, amount)
transactions_data = [
(1, 1000, 188.23),
(2, 1001, 0.50), # Small transaction
(2, 1002, 600.00), # Large transaction - ALERT!
(3, 1003, 42.00),
(1, 1004, 0.89), # Small transaction
(1, 1005, 300.00), # Not large enough - no alert
(4, 1006, 0.10), # Small transaction
(4, 1007, 520.00), # Large transaction - ALERT!
(3, 1008, 0.75), # Small transaction
(3, 1009, 800.00), # Large transaction - ALERT!
]
transactions = env.from_collection(
transactions_data,
type_info=Types.TUPLE([Types.LONG(), Types.LONG(), Types.DOUBLE()])
)
alerts = transactions \
.key_by(lambda t: t[0]) \
.process(FraudDetector())
alerts.print()
env.execute("Fraud Detection")
if __name__ == '__main__':
fraud_detection()
逐步理解代码
主类定义应用的数据流;FraudDetector 定义识别交易模式的业务逻辑。
执行环境
首先取得 StreamExecutionEnvironment。通过执行环境设置作业属性、创建数据源,并触发作业执行。
Java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Python
env = StreamExecutionEnvironment.get_execution_environment()
创建数据源
source 把 Apache Kafka、RabbitMQ、Apache Pulsar 等外部系统的数据接入 Flink 作业。
Java
本示例使用 TransactionSource,它封装 DataGeneratorSource,持续生成无界的示例交易。每条交易包含账户 ID accountId、发生时间戳 timestamp 和金额 amount。
DataStream<Transaction> transactions = env
.fromSource(
TransactionSource.unbounded(),
WatermarkStrategy.noWatermarks(),
"transactions");
fromSource 的三个参数分别是 source、水位线策略和用于调试的名称。这里使用处理时间,不使用事件时间水位线,因此传入 noWatermarks()。
Python
Python 示例使用有限的样本集合。每笔交易是一个包含账户 ID、时间戳和金额的元组:
transactions_data = [
(1, 1000, 188.23),
(2, 1001, 0.50), # Small transaction
(2, 1002, 600.00), # Large transaction - ALERT!
# ...more transactions
]
transactions = env.from_collection(
transactions_data,
type_info=Types.TUPLE([Types.LONG(), Types.LONG(), Types.DOUBLE()])
)
from_collection 便于教学和测试。生产系统通常使用 Kafka 等 source 连接器。
按账户分区并检测模式
transactions 流包含许多用户的交易,可由多个任务并行处理。检测规则以账户为单位,因此同一账户的交易必须到达同一个检测算子的并行任务。
keyBy 按键分区,确保同一个键的记录由同一物理任务处理。process() 添加算子,逐条应用处理函数。通常说紧跟 keyBy 的算子运行在 keyed context 中;这里就是 FraudDetector。
Java
DataStream<Alert> alerts = transactions
.keyBy(Transaction::getAccountId)
.process(new FraudDetector())
.name("fraud-detector");
Python
alerts = transactions \
.key_by(lambda t: t[0]) \
.process(FraudDetector())
输出结果
sink 把 DataStream 写入外部系统,例如 Kafka、Cassandra 或 AWS Kinesis。
Java
AlertSink 仅以 INFO 级别记录每条 Alert,便于观察结果,没有把它们写入持久化存储。
alerts.addSink(new AlertSink());
Python
print() 把告警输出到控制台:
alerts.print()
检测函数
检测器是一个 KeyedProcessFunction,每笔交易都会调用其 processElement,Python 中为 process_element。下面先展示最简单的版本:每笔交易都产生告警,显然比较保守。随后再加入真正的判断逻辑。
Java
public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {
private static final double SMALL_AMOUNT = 1.00;
private static final double LARGE_AMOUNT = 500.00;
@Override
public void processElement(
Transaction transaction,
Context context,
Collector<Alert> collector) throws Exception {
Alert alert = new Alert();
alert.setId(transaction.getAccountId());
collector.collect(alert);
}
}
Python
class FraudDetector(KeyedProcessFunction):
SMALL_AMOUNT = 1.00
LARGE_AMOUNT = 500.00
def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
account_id = transaction[0]
yield f"Alert{{id={account_id}}}"
实现业务规则
第一版规则是:同一账户先发生一笔小额交易,紧接着发生大额交易,就发出告警。这里把小额定义为低于 1.00,大额定义为高于 500,等于阈值不属于对应类别。
原教程的 交易顺序示意图说明两种情况。关键关系如下,金额均为教学数据:
| 交易顺序 | 模式 | 结果 |
|---|---|---|
| 3 → 4 | 0.09 后紧接 510 | 触发告警 |
| 7 → 8 → 9 | 0.02 后有一笔中间交易,再出现大额交易 | 中间交易打断相邻模式,不告警 |
检测器必须跨事件记住信息:只有前一笔交易是小额时,当前的大额交易才符合规则。这需要 状态,也是采用 KeyedProcessFunction 的原因。它允许细致控制状态和时间,便于继续增加复杂要求。
最直接的思路是在遇到小额交易时设置布尔标志,下一笔大额交易到来时检查账户对应的标志。但不能简单把它放进 FraudDetector 的普通成员变量:同一算子实例会处理多个账户,账户 A 设置的标志可能导致账户 B 被误报。
即使使用 Map 为各账户保存标志,普通成员数据也不会自动成为 Flink 的托管状态。故障后重启会丢失这些信息,可能漏报。Flink 提供与普通变量相近的状态原语,让框架管理状态及其恢复。
建立 ValueState
最基础的 ValueState属于 keyed state,只能在按键上下文中使用,例如紧跟 keyBy 的算子。访问状态时,框架自动按当前记录的键限定范围。这里的键是账户 ID,因此每个账户都有独立状态。
创建状态时使用 ValueStateDescriptor,说明框架应如何管理这个变量。在开始处理数据前,通过 open() 注册状态:
Java
public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {
private static final long serialVersionUID = 1L;
private transient ValueState<Boolean> flagState;
@Override
public void open(OpenContext openContext) {
ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
"flag",
Types.BOOLEAN);
flagState = getRuntimeContext().getState(flagDescriptor);
}
Python
class FraudDetector(KeyedProcessFunction):
def __init__(self):
self.flag_state = None
def open(self, runtime_context: RuntimeContext):
descriptor = ValueStateDescriptor("flag", Types.BOOLEAN())
self.flag_state = runtime_context.get_state(descriptor)
ValueState 是包装类型,类似 Java 的 AtomicReference 或 AtomicLong。update 设置值,value 读取值,clear 清空值。一个键尚无状态或被清空后,读取结果为 Java 的 null,或 Python 的 None。
直接修改 value 返回的对象,不保证被状态系统识别;变更必须通过 update 提交。状态的保存与恢复由 Flink 管理,但示例本身没有启用 checkpoint;故障恢复还需要实际配置检查点或使用保存点,不能仅凭声明 ValueState 就保证故障后恢复最近状态。
用标志追踪相邻交易
Java
@Override
public void processElement(
Transaction transaction,
Context context,
Collector<Alert> collector) throws Exception {
// Get the current state for the current key
Boolean lastTransactionWasSmall = flagState.value();
// Check if the flag is set
if (lastTransactionWasSmall != null) {
if (transaction.getAmount() > LARGE_AMOUNT) {
// Output an alert downstream
Alert alert = new Alert();
alert.setId(transaction.getAccountId());
collector.collect(alert);
}
// Clean up our state
flagState.clear();
}
if (transaction.getAmount() < SMALL_AMOUNT) {
// Set the flag to true
flagState.update(true);
}
}
Python
def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
account_id = transaction[0]
amount = transaction[2]
# Get the current state for the current key
last_transaction_was_small = self.flag_state.value()
# Check if the flag is set
if last_transaction_was_small is not None:
if amount > self.LARGE_AMOUNT:
# Output an alert downstream
yield f"Alert{{id={account_id}}}"
# Clean up our state
self.flag_state.clear()
if amount < self.SMALL_AMOUNT:
# Set the flag to true
self.flag_state.update(True)
每笔交易先读取该账户的标志。ValueState 的范围始终是当前账户。标志非空意味着该账户上一笔交易是小额;如果当前金额是大额,就输出告警。
完成检查后,清除前一笔留下的标志:若已告警,模式结束;若未告警,相邻模式已经被当前交易打断。最后检查当前交易是否为小额,若是则设置标志,留给下一笔交易检查。
ValueState<Boolean> 可处于未设置、真、假三种状态,因为 ValueState 可以为空。这里仅使用“未设置”和“真”,通过是否为空判断标志,而不会写入假值。
完整实现
Java
下面是带托管状态的完整 FraudDetector:
import org.apache.flink.api.common.functions.OpenContext;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.util.Collector;
import org.apache.flink.walkthrough.common.entity.Alert;
import org.apache.flink.walkthrough.common.entity.Transaction;
public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {
private static final long serialVersionUID = 1L;
private static final double SMALL_AMOUNT = 1.00;
private static final double LARGE_AMOUNT = 500.00;
private transient ValueState<Boolean> flagState;
@Override
public void open(OpenContext openContext) {
ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
"flag",
Types.BOOLEAN);
flagState = getRuntimeContext().getState(flagDescriptor);
}
@Override
public void processElement(
Transaction transaction,
Context context,
Collector<Alert> collector) throws Exception {
// Get the current state for the current key
Boolean lastTransactionWasSmall = flagState.value();
// Check if the flag is set
if (lastTransactionWasSmall != null) {
if (transaction.getAmount() > LARGE_AMOUNT) {
// Output an alert downstream
Alert alert = new Alert();
alert.setId(transaction.getAccountId());
collector.collect(alert);
}
// Clean up our state
flagState.clear();
}
if (transaction.getAmount() < SMALL_AMOUNT) {
// Set the flag to true
flagState.update(true);
}
}
}
Python
下面是完整的 Python 程序:
from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.functions import KeyedProcessFunction, RuntimeContext
from pyflink.datastream.state import ValueStateDescriptor
class FraudDetector(KeyedProcessFunction):
SMALL_AMOUNT = 1.00
LARGE_AMOUNT = 500.00
def __init__(self):
self.flag_state = None
def open(self, runtime_context: RuntimeContext):
descriptor = ValueStateDescriptor("flag", Types.BOOLEAN())
self.flag_state = runtime_context.get_state(descriptor)
def process_element(self, transaction, ctx: 'KeyedProcessFunction.Context'):
# transaction is a tuple: (account_id, timestamp, amount)
account_id = transaction[0]
amount = transaction[2]
# Get the current state for the current key
last_transaction_was_small = self.flag_state.value()
# Check if the flag is set
if last_transaction_was_small is not None:
if amount > self.LARGE_AMOUNT:
# Output an alert downstream
yield f"Alert{{id={account_id}}}"
# Clean up our state
self.flag_state.clear()
if amount < self.SMALL_AMOUNT:
# Set the flag to true
self.flag_state.update(True)
def fraud_detection():
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)
# Sample transaction data: (account_id, timestamp, amount)
transactions_data = [
(1, 1000, 188.23),
(2, 1001, 0.50), # Small transaction
(2, 1002, 600.00), # Large transaction - ALERT!
(3, 1003, 42.00),
(1, 1004, 0.89), # Small transaction
(1, 1005, 300.00), # Not large enough - no alert
(4, 1006, 0.10), # Small transaction
(4, 1007, 520.00), # Large transaction - ALERT!
(3, 1008, 0.75), # Small transaction
(3, 1009, 800.00), # Large transaction - ALERT!
]
transactions = env.from_collection(
transactions_data,
type_info=Types.TUPLE([Types.LONG(), Types.LONG(), Types.DOUBLE()])
)
alerts = transactions \
.key_by(lambda t: t[0]) \
.process(FraudDetector())
alerts.print()
env.execute("Fraud Detection")
if __name__ == '__main__':
fraud_detection()
运行应用
到这里,程序的数据源、按账户处理、状态与输出已连接起来。Java 输入是无界的,会持续处理,直到手动停止;Python 版本处理有限样本集合。
Java
在 IDE 中运行 FraudDetectionJob,可以观察流处理结果。原文展示的日志示例如下,时间戳属于原文示例,不是本机运行记录:
2024-01-01 14:22:06,220 INFO org.apache.flink.walkthrough.common.sink.AlertSink - Alert{id=3}
2024-01-01 14:22:11,383 INFO org.apache.flink.walkthrough.common.sink.AlertSink - Alert{id=3}
2024-01-01 14:22:16,551 INFO org.apache.flink.walkthrough.common.sink.AlertSink - Alert{id=3}
Python
可在命令行运行:
$ python fraud_detection.py
原文展示的输出示例如下:
Alert{id=2}
Alert{id=4}
Alert{id=3}
该命令在本地 mini cluster 中构建并运行 PyFlink 程序,也可把作业提交到远程集群,详见 作业提交示例。本文仅核对代码和控制流,没有执行 Flink 或这些命令。
下一步
给规则增加时间条件
当前实现只检测同一账户的“小额后紧接大额”模式,没有时间限制。实际检测往往还需要时间条件。要学习定时器,例如仅检测一分钟内出现的两笔交易,可继续阅读 事件驱动应用。这里的规则是教学模型,不能直接代表已经验证的真实反欺诈系统。
深入 DataStream
- DataStream API 概览:完整 API 参考。
- Process Functions:深入了解 KeyedProcessFunction 等处理函数。
- 状态与容错:了解托管状态和 checkpoint。
其他教程
- Flink SQL 教程:不编写程序的交互式 SQL。
- Table API 教程:用 Table API 构建流处理管道。
- Flink Operations Playground:学习集群运维。











暂无评论内容