Flink RocksDB StateBackend并行度:算子如何访问RocksDB表?
Flink作业技术疑问解答
作业代码
package spendreport; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.walkthrough.common.sink.AlertSink; import org.apache.flink.walkthrough.common.entity.Alert; import org.apache.flink.walkthrough.common.entity.Transaction; import org.apache.flink.walkthrough.common.source.TransactionSource; public class FraudDetectionJob { private String checkpointsDir = "file://checkpoints/"; private String rocksDBStateDir = "file://state/rocksdb/"; public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(3); EnvironmentSettings tableSettings = EnvironmentSettings .newInstance() .useBlinkPlanner() .build(); StreamTableEnvironment tableEnv = StreamTableEnvironment .create(env, tableSettings); tableEnv.executeSql("CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...)"); env.enableCheckpointing(5000); env.checkpointConfig.minPauseBetweenCheckpoints = 100; env.checkpointConfig.setCheckpointStorage(checkpointsDir); StateBackend stateBackend = new EmbeddedRocksDBStateBackend(); stateBackend.setDbStoragePath(rocksDBStateDir); env.setStateBackend(stateBackend); DataStream<Transaction> transactions = env .addSource(new TransactionSource()) .name("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"); } }
技术疑问与解答
疑问1:作业中的每个算子实例是否会访问同一个RocksDB表Orders进行读写操作?
不会。首先要明确:你通过CREATE TABLE定义的Orders表,其存储介质由WITH子句的配置决定,和作业配置的EmbeddedRocksDBStateBackend完全是两回事——后者是用来存储算子状态(比如FraudDetector里维护的业务状态)的,并非用来存储Orders表的数据。
即便Orders表恰好采用RocksDB作为存储,Flink的分布式架构也会让每个算子并行实例处理对应分区的数据,各自维护独立的RocksDB分片,绝不会共用同一个RocksDB表进行读写操作。
疑问2:作业级并行度是否意味着每个作业拥有独立的RocksDB实例?因为RocksDB是基于TaskManager的,无法保证所有作业都运行在同一个TaskManager上。
是的,每个作业都会拥有独立的RocksDB实例,原因如下:
- Flink严格隔离不同作业的状态,即便多个作业运行在同一个TaskManager上,它们的RocksDB数据目录、资源都是完全隔离的,各自拥有独立的实例。
- 作业级并行度控制的是单个作业内算子的并行实例数量,每个并行实例的RocksDB状态(基于
EmbeddedRocksDBStateBackend)会存储在配置的rocksDBStateDir下的作业专属路径中,不同作业之间不会互相干扰。 - 哪怕作业的并行实例分布在不同TaskManager上,每个TaskManager上的作业实例都会有自己的RocksDB实例,共同构成作业的分布式状态存储,但作业之间依然保持隔离。
内容的提问来源于stack exchange,提问作者overexchange
相关产品推荐
相关产品推荐

