The source notes for this track are currently maintained in Chinese.
预计阅读时间: 90 分钟 前置阅读: doc-03(Checkpoint) 下一次阅读: doc-05(Window)
1. 状态分类与内部表示
1.1 三种状态类型
| 类型 | 作用域 | 存储后端 | 典型场景 | 内部结构 |
|---|---|---|---|---|
| KeyedState | keyBy() 后的每个 Key | StateBackend (Heap/RocksDB) | GROUP BY, 窗口聚合, 计数 | KeyGroup → {Key → Value} |
| OperatorState | 整个算子实例 | StateBackend | Kafka Source offset | List<StateHandle> |
| BroadcastState | 广播到所有 Subtask | StateBackend | 动态配置分发 | Key → Value (每个 Subtask 全量) |
1.2 KeyedState 的物理组织 — KeyGroup
Key 不直接映射到 Subtask,而是经过两层映射:
Key → hash(key) % maxParallelism → KeyGroup → KeyGroupRangeAssignment → Subtask
例如 maxParallelism=128, parallelism=4:
KeyGroup 分配:
Subtask 0: KeyGroup [0, 31]
Subtask 1: KeyGroup [32, 63]
Subtask 2: KeyGroup [64, 95]
Subtask 3: KeyGroup [96, 127]
Key "user_123" → hash=42 → KeyGroup=42 → Subtask 1 (KeyGroup [32,63])
// KeyGroupRangeAssignment.java — 核心分配算法
public static int computeKeyGroupForKeyHash(int keyHash, int maxParallelism) {
return MathUtils.murmurHash(keyHash) % maxParallelism;
}
public static int computeOperatorIndexForKeyGroup(
int maxParallelism, int parallelism, int keyGroup) {
// 将 KeyGroup 均匀分配到 Subtask
return keyGroup * parallelism / maxParallelism;
}
为什么需要 KeyGroup 这一层?
- KeyGroup 数量(maxParallelism)创建后固定不变
- 当并行度改变时,只需重新分配 KeyGroup→Subtask 映射
- Key→KeyGroup 映射永远不变 → Rescale 时状态可以恢复
1.3 五种 KeyedState 类型
// KeyedStateBackend.java — 状态类型
public interface KeyedStateBackend<K> extends KeyedStateFactory {
// 1. ValueState<T>: 单值状态
<T> ValueState<T> getPartitionedState(
ValueStateDescriptor<T> descriptor);
// 2. ListState<T>: 列表状态
<T> ListState<T> getPartitionedState(
ListStateDescriptor<T> descriptor);
// 3. MapState<K,V>: Map 状态 (支持迭代)
<T, UK, UV> MapState<UK, UV> getPartitionedState(
MapStateDescriptor<UK, UV> descriptor);
// 4. ReducingState<T>: 归约状态 (自动 reduce)
<T> ReducingState<T> getPartitionedState(
ReducingStateDescriptor<T> descriptor);
// 5. AggregatingState<IN, OUT>: 聚合状态
<IN, ACC, OUT> AggregatingState<IN, OUT> getPartitionedState(
AggregatingStateDescriptor<IN, ACC, OUT> descriptor);
}
ValueState vs MapState 的选择:
| 维度 | ValueState | MapState |
|---|---|---|
| 存储 | 单值 | Key-Value 集合 |
| State TTL | 支持 | 每个 Entry 独立 TTL |
| 迭代 | 不支持 | 支持 keys()/values()/entries() 迭代 |
| RocksDB 存储 | 单 CF Key-Value | 多个 CF Key-Value (带前缀) |
| 适用 | 简单计数、累加 | 窗口内元素、GROUP BY 多个 group |
2. StateBackend 深度对比
2.1 三种 Backend 对比
| 维度 | HashMapStateBackend | EmbeddedRocksDBStateBackend | ChangelogStateBackend |
|---|---|---|---|
| 存储位置 | Java Heap | RocksDB (本地磁盘) | DFS Changelog+本地 Materialized |
| 状态大小上限 | ~100MB (受 GC 限制) | TB 级 (受本地磁盘限制) | 中(受 DFS 限制) |
| 访问速度 | 极快(直接堆访问) | 快(有序列化开销,~1-5μs) | 快(本地读)+ DFS 写 |
| 序列化 | 仅在 Checkpoint 时 | 每次读写都需要 | 写入时序列化到 DFS |
| Checkpoint 格式 | 全量 (Java 对象 → 字节) | 增量 (SST 文件 diff) | Changelog 日志 |
| 恢复速度 | 快(纯内存加载) | 中(需下载 SST + 重放 WAL) | 极快(Apply Changelog) |
| GC 影响 | 大(状态在 Heap 中) | 无(状态在 Off-Heap/磁盘) | 中 |
| 适用场景 | 小状态 + 低延迟 | 大状态 + 稳定吞吐 | 需要快速恢复的场景 |
2.2 StateBackend 接口
// StateBackend.java — 创建 KeyedStateBackend 工厂
public interface StateBackend extends java.io.Serializable {
// 创建 KeyedStateBackend(keyBy 后的状态)
<K> KeyedStateBackend<K> createKeyedStateBackend(
Environment env, JobID jobID, String operatorIdentifier,
TypeSerializer<K> keySerializer, int numberOfKeyGroups,
KeyGroupRange keyGroupRange, TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider, MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles, // 恢复时的状态
CloseableRegistry cancelStreamRegistry,
double managedMemoryFraction) throws Exception;
// 创建 OperatorStateBackend(Kafka Source offset 等)
OperatorStateBackend createOperatorStateBackend(
Environment env, String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) throws Exception;
}
3. EmbeddedRocksDBStateBackend 内部原理
3.1 架构总览
TaskManager JVM
│
├── RocksDBKeyedStateBackend
│ │
│ ├── ColumnFamily per State:
│ │ ├── CF "windowState" (MapState 的每个 key-value 用一个 RocksDB key)
│ │ ├── CF "countState" (ValueState 单值)
│ │ └── CF "default"
│ │
│ ├── WriteBatch (批量写入,减少 RocksDB 调用)
│ │ └── 积累多次 state update → 批量 flush
│ │
│ ├── RocksDB 实例(C++ JNI 调用)
│ │ ├── MemTable (write buffer) → Immutable MemTable → flush
│ │ ├── SST Files: L0 (overlapping) → L1 → ... → L6 (non-overlapping)
│ │ ├── Block Cache (读缓存,在 Managed Memory 中)
│ │ └── WAL (write-ahead log, 进程崩溃恢复)
│ │
│ └── Managed Memory 分配:
│ ├── writeBufferSize * maxWriteBufferNumber → MemTable
│ ├── blockCacheSize → Block Cache
│ └── 总 managed memory 在 taskmanager.memory.managed.size 中配置
3.2 RocksDB 中的 Key 编码
RocksDB 是 KV 存储,Flink 通过组合编码将 KeyedState 存储在 RocksDB 中:
RocksDB Key 格式:
┌──────────────────────────────────────────────────┐
│ KeyGroup(2B) │ Key(Variable) │ Namespace(Variable) │
└──────────────────────────────────────────────────┘
例如:
KeyGroup=42, Key="user_123", Namespace=Window[10:00,10:01)
→ RocksDB Key = [0x002A][user_123][window-10:00-10:01]
为什么要加 KeyGroup 前缀?
- 恢复时按 KeyGroup 扫描(不需要全量扫描)
- KeyGroup 是恢复和 Rescale 的最小单位
3.3 RocksDB 写路径
// RocksDBWriteBatchWrapper.java — 批量写入
public class RocksDBWriteBatchWrapper implements AutoCloseable {
private final WriteBatch writeBatch;
private final WriteOptions writeOptions; // disable WAL 等配置
public void put(ColumnFamilyHandle cf, byte[] key, byte[] value) throws RocksDBException {
if (writeBatch.count() >= batchSize || sizeInBytes >= batchSizeInBytes) {
flush(); // 批量满 → 写入 RocksDB
}
writeBatch.put(cf, key, value);
sizeInBytes += key.length + value.length;
}
private void flush() throws RocksDBException {
rocksDB.write(writeOptions, writeBatch);
writeBatch.clear();
}
}
写入优化:
- WriteBatch 批量积累:默认 batchSize=2MB,减少 RocksDB write() 调用
- 禁用 WAL(可选):Checkpoint 提供容错,RocksDB WAL 可禁用 → 写性能大幅提升
- MemTable 快速路径:大部分写入直接写 MemTable(内存),异步 flush 到磁盘
3.4 RocksDB 读路径
// RocksDBValueState.java — 读取 ValueState
public V value() throws IOException {
byte[] key = serializeKeyAndNamespace(currentKey, namespace);
byte[] valueBytes = backend.db.get(
currentKeyGroup, // KeyGroup 范围查询
columnFamily, // State 对应的 CF
key
);
if (valueBytes == null) {
return defaultValue;
}
return serializer.deserialize(valueBytes);
}
读的层次(RocksDB 内部):
MemTable → Immutable MemTable → Block Cache → SST Files (L0→L6)
↑
Managed Memory 中的 Block Cache
缓存频繁读取的 SST Block
RocksDB 的序列化开销:每次 get() 和 put() 都需要将 Java 对象序列化为字节数组。这是 RocksDB Backend 的主要性能开销来源。对于访问非常频繁的小状态,HashMapStateBackend 没有这个开销。
4. HashMapStateBackend — 纯堆内存状态
4.1 内部结构
// HeapKeyedStateBackend.java — 核心存储
public class HeapKeyedStateBackend<K> extends AbstractKeyedStateBackend<K> {
// 存储: Map<StateName, StateTable>
// StateTable = Map<KeyGroup, Map<K, State>>
private final Map<String, StateTable<K, ?, ?>> registeredKVStates;
// 每个 StateDescriptor 对应一个 StateTable
// StateTable 内部是一组 CopyOnWriteStateMap
// (每个 KeyGroup 一个 CopyOnWriteStateMap)
}
4.2 CopyOnWriteStateMap — 并发安全的状态读写
HeapKeyedStateBackend 的并发模型:
每个 KeyGroup 有独立的 CopyOnWriteStateMap<K, N, S>
- 读写分离: 读取无锁,写入时拷贝整个 Map
- 适用于小状态的场景(小于 100MB)
- GC 压力: 每次 Checkpoint 生成的快照是堆上的新对象
CopyOnWriteStateMap 结构:
StateMapEntry[] table (hash table)
└── StateMapEntry { key, namespace, state, next(in-chain) }
4.3 Checkpoint 快照
// HeapKeyedStateBackend.snapshot() — 全量快照
public void snapshot(long checkpointId, long timestamp,
CheckpointStreamFactory streamFactory,
CheckpointOptions checkpointOptions) throws Exception {
// 1. 遍历所有 State → 序列化为字节数组
for (StateTable stateTable : registeredKVStates.values()) {
// 2. 按 KeyGroup 分组序列化
for (int keyGroup : keyGroupRange) {
// 3. 序列化 KeyGroup 内的所有 Key→State
DataOutputView outputView = ...;
stateTable.writeStateInKeyGroup(outputView, keyGroup);
}
}
// 4. 写入 DFS (HDFS/S3)
// 全量写入 → 状态越大,Checkpoint 越慢
}
HashMapStateBackend 的 Checkpoint 限制:
- 全量 Checkpoint,每次都要序列化和上传全部状态
- 状态 > 1GB 时,Checkpoint 耗时可能超过 1 分钟
- 状态 > 100MB 时,GC 压力显著增加
- 推荐:状态 < 100MB → HashMapStateBackend;> 100MB → RocksDB
5. ChangelogStateBackend — 基于日志的状态
5.1 设计思路
ChangelogStateBackend 将状态变更记录到分布式日志中:
┌─────────────────────────────────────────────┐
│ ChangelogStateBackend │
│ │
│ 状态写入: state.update(value) │
│ │ │
│ ├──→ 写入本地 Materialized Store (快) │
│ └──→ 追加 Changelog Log (DFS) │
│ │
│ 恢复时: │
│ └──→ 从 DFS 读取 Changelog │
│ → Apply 到本地 Store │
│ → 恢复完成(无需下载全量状态) │
└─────────────────────────────────────────────┘
5.2 Changelog vs 传统 Checkpoint
| 维度 | 传统 Checkpoint | Changelog Backend |
|---|---|---|
| 恢复时间 | 下载所有 SST → 重放 WAL | Apply Changelog 日志 |
| 恢复数据量 | O(状态大小) 如 100GB | O(Changelog 增量) 如 100MB |
| 写入开销 | 只在 Checkpoint 时写 DFS | 每次状态更新都写 DFS(可能 durably 或 periodically) |
| 一致性 | Checkpoint 间隔内无持久化 | Periodic materialize + continuous log |
6. State TTL — 状态过期机制
6.1 TTL 配置
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24)) // TTL 24 小时
.setUpdateType(UpdateType.OnCreateAndWrite) // 写时更新时间戳
.setStateVisibility(StateVisibility.NeverReturnExpired) // 不返回过期值
.setTtlTimeCharacteristic(TtlTimeCharacteristic.ProcessingTime) // 处理时间
.cleanupFullSnapshot() // Checkpoint 时全量清理
.cleanupIncrementally(10, false) // 每次访问时增量清理
.cleanupInRocksdbCompactFilter() // RocksDB Compaction 时清理
.build();
6.2 TTL 清理策略详解
| 策略 | 机制 | 开销 | 适用 |
|---|---|---|---|
| cleanupFullSnapshot | Checkpoint 时扫描所有状态,丢弃过期项 | Checkpoint 开销增加 | 默认策略,最安全 |
| cleanupIncrementally | 每次 state.value() 访问时检查并清理过期项 | 每次访问有微小开销 | 访问频繁但状态量大的场景 |
| cleanupInRocksdbCompactFilter | RocksDB Compaction 时通过 CompactionFilter 清理 | 零读取开销 | 大状态 + RocksDB + 不要求精确到期 |
| cleanupFullSnapshot + cleanupIncrementally | 组合 | 均衡 | 生产推荐 |
6.3 TTL 的内部实现
// TtlValueStateWrapper.java — TTL 包装器
class TtlValueStateWrapper<T> extends AbstractTtlState<T> implements ValueState<T> {
// 每次写时记录 timestamp
@Override
public void update(T value) {
long timestamp = timeProvider.currentTimestamp(); // 当前处理时间/事件时间
wrappedState.update(TtlValue.wrap(value, timestamp));
}
// 每次读时检查 TTL
@Override
public T value() {
TtlValue<T> ttlValue = wrappedState.value();
if (ttlValue == null) return null;
if (expired(ttlValue)) { // timestamp + ttl < currentTime
// 根据 visibility 决定是否返回
if (visibility == NeverReturnExpired) {
return null; // 或 defaultValue
}
}
return ttlValue.getValue();
}
}
TTL 的关键时序问题:
- ProcessingTime: 精确(系统时间定长),但重启后时间重置
- EventTime: 受 Watermark 推进影响,如果 Watermark 停滞,TTL 不生效 → 状态可能堆积
setUpdateType(OnReadAndWrite): 每次读也更新 TTL,防止"热"数据意外过期
7. Managed Memory 与状态的关系
7.1 Managed Memory 在 RocksDB Backend 中的作用
taskmanager.memory.managed.size (默认: flink.size * 0.4)
┌─────────────────────────────────────────────────┐
│ Managed Memory = 1GB │
├─────────────────────────────────────────────────┤
│ RocksDB Write Buffer (MemTable) = 512MB │
│ RocksDB Block Cache (Read Cache) = 256MB │
│ Sort-Merge Join 缓冲区 = 128MB │
│ Hash Join Build Side = 96MB │
│ Python UDF 内存 = 32MB │
└─────────────────────────────────────────────────┘
RocksDB 核心配置:
| ROCKSDD 配置 | 默认 | 说明 |
|---|---|---|
write_buffer_size | 64MB | 单个 MemTable 的大小 |
max_write_buffer_number | 2 | MemTable 的数量(1 active + 1 immutable) |
block_cache_size | 8MB (per CF) | 读缓存的 Block Cache 大小 |
target_file_size_base | 64MB | L1 层 SST 文件目标大小 |
max_bytes_for_level_base | 256MB | L1 层总大小上限 |
7.2 MemoryManager 源码
// MemoryManager.java — Managed Memory 的池化管理
public class MemoryManager {
private final long memorySize; // 总内存大小
private final int pageSize; // 内存页大小 (32KB)
private final int numberOfPages; // 总页数
private final ArrayDeque<MemorySegment> freeSegments; // 空闲段池
// 分配内存段
public List<MemorySegment> allocatePages(
Object owner, int numPages) throws MemoryAllocationException {
synchronized (lock) {
if (freeSegments.size() < numPages) {
throw new MemoryAllocationException(...);
}
List<MemorySegment> allocated = new ArrayList<>(numPages);
for (int i = 0; i < numPages; i++) {
allocated.add(freeSegments.poll());
}
allocatedSegments.put(owner, allocated); // 记账
return allocated;
}
}
// 释放内存段
public void release(List<MemorySegment> segments) {
synchronized (lock) {
freeSegments.addAll(segments);
// 更新记账
}
}
}
8. RocksDB 增量 Checkpoint 深度流程
8.1 完整流程
RocksDBStateBackend.snapshotState(checkpointId=42)
│
├── Step 1: flushMemTable()
│ └── 将 MemTable 刷到 SST 文件 → 生成新的 SST 文件
│
├── Step 2: rocksDB.getLiveFiles()
│ └── 获取当前所有 SST 文件的列表
│ eg: [SST_01, SST_02, SST_03, SST_04, SST_05]
│
├── Step 3: chooseSstFilesForUpload()
│ └── 对比 lastCheckpoint (Chk#41) 的 SST 列表
│ Chk#41 SST files: [SST_01, SST_02, SST_03]
│ Current SST files: [SST_01, SST_02, SST_03, SST_04, SST_05]
│ 新增: [SST_04, SST_05] ← 只需上传这两个
│ 删除: [] ← compaction 没有删除文件
│
├── Step 4: uploadFilesToCheckpointFs()
│ └── 上传 SST_04, SST_05 到 DFS
│ 异步执行: 不影响数据处理
│
├── Step 5: SharedStateRegistry.registerReference()
│ └── 注册 SST_01, SST_02, SST_03 被 Chk#42 引用
│ refCount[SST_01] = 2 (Chk#41 + Chk#42)
│ refCount[SST_02] = 2
│ refCount[SST_03] = 2
│ refCount[SST_04] = 1 (新)
│ refCount[SST_05] = 1 (新)
│
└── Step 6: 构造 IncrementalRemoteKeyedStateHandle
→ 返回给 CheckpointCoordinator
→ 包含: uploadedHandles + sharedStateRegistry references
8.2 SharedStateRegistry 的引用计数 GC
SharedStateRegistry 维护:
Map<SstFileKey, Set<CheckpointId>> // 每个 SST 文件被哪些 Checkpoint 引用
GC 逻辑:
Chk#41 过期 (被新 Checkpoint 淘汰或达到 maxCheckpoints):
→ 遍历 Chk#41 的 SST 文件
→ refCount[SST_01]-- (2→1)
→ refCount[SST_02]-- (2→1)
→ refCount[SST_03]-- (2→1)
→ SST_01, SST_02, SST_03 仍被 Chk#42 引用 → 保留
Chk#42 过期:
→ refCount[SST_01]-- (1→0) → 删除 DFS 上的 SST_01
→ refCount[SST_02]-- (1→0) → 删除 DFS 上的 SST_02
→ ...
9. 状态序列化
9.1 TypeSerializer 体系
TypeSerializer<T> (Flink 的序列化抽象)
├── BasicTypeSerializer (Int, Long, String, Double, Boolean...)
├── PojoSerializer (POJO 类型)
├── KryoSerializer (通用后备方案)
├── RowSerializer (Row 类型, Table API 内部)
├── TupleSerializer (Tuple 类型)
└── CompositeTypeSerializer (自定义复合类型)
序列化性能对比:
Kryo (通用) ~100 ns/object (慢, 但通用)
PojoSerializer ~50 ns/object (中)
BasicType ~5 ns/object (快, 仅基本类型)
Row ~20 ns/object (快, Table API 优化)
9.2 RocksDB 中的序列化路径
State 写入:
Java Object → TypeSerializer.serialize() → byte[] → RocksDB.put(keyBytes, valueBytes)
State 读取:
RocksDB.get(key) → byte[] → TypeSerializer.deserialize() → Java Object
Key 编码:
(KeyGroup, Key, Namespace) → CompositeKeySerialization.serializeKeyGroupNamespaceKey()
特殊优化: 当 Key 和 Namespace 是固定长度类型时(如 Long),使用高效的编码路径
减少序列化开销的技巧:
- 使用 POJO 而非 Kryo(
env.getConfig().enableForceKryo()→ 关闭) - 避免大对象作为 State Value(拆成多个小 ValueState)
- RocksDB 的 Block Cache 可以缓存序列化后的字节,减少 deserialize 次数
10. 源码导航(完整版)
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-runtime/.../state/StateBackend.java | createKeyedStateBackend() | StateBackend 接口 |
flink-runtime/.../state/KeyedStateBackend.java | setCurrentKey(), getPartitionedState() | KeyedState Backend |
flink-runtime/.../state/heap/HeapKeyedStateBackend.java | snapshot(), notifyCheckpointComplete() | 堆状态后端 |
flink-runtime/.../state/heap/CopyOnWriteStateMap.java | get(), put(), remove() | 堆状态存储结构 |
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBKeyedStateBackend.java | snapshot() | RocksDB 状态后端 |
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBWriteBatchWrapper.java | put(), flush() | RocksDB 批量写入 |
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBIncrementalCheckpointUtils.java | chooseSstFilesForUpload() | 增量 SST 选择 |
flink-runtime/.../state/ttl/TtlStateFactory.java | createValueState(), ... | TTL 状态包装器 |
flink-runtime/.../state/ttl/TtlValueState.java | value(), update() | TTL ValueState 实现 |
flink-runtime/.../memory/MemoryManager.java | allocatePages(), release() | Managed Memory 管理 |
flink-runtime/.../state/KeyGroupRangeAssignment.java | computeKeyGroupForHash(), computeOperatorIndexForKeyGroup() | KeyGroup 分配 |
flink-runtime/.../state/SharedStateRegistry.java | registerReference(), unregisterUnusedState() | SST 文件引用计数 |
11. 常见问题 / 面试题
Q1: RocksDB vs HashMap StateBackend 的选择标准?
A: 以 100MB 为界。以下是详细决策树:
- 状态 < 100MB → HashMapStateBackend(零序列化开销,低延迟)
- 状态 > 100MB → EmbeddedRocksDBStateBackend(无 GC 压力,增量 Checkpoint)
- 状态巨大但需要极快恢复 → ChangelogStateBackend
- 不确定 → 从 RocksDB 开始,因为它不增加 GC 压力且可增量 Checkpoint
Q2: RocksDB 的 WAL 和 Flink Checkpoint 的关系是什么?
A: 它们互补但不重复:
- RocksDB WAL:进程级崩溃恢复(TM 重启但本地磁盘还在),恢复 ~秒级
- Flink Checkpoint:全局容错(TM 挂掉、磁盘损坏、集群故障),恢复 ~分钟级
- 默认情况下 Flink Checkpoint 前会 flush MemTable,因此 WAL 在 Checkpoint 间的作用有限
- 如果禁用 WAL(
state.backend.rocksdb.wal.disabled: true),依赖 Checkpoint 恢复
Q3: 状态大小如何监控?
A: 关键 Metrics:
State Size:每个算子的状态大小(RocksDB 的estimate-num-keys,size-all-mem-tables,block-cache-usage)Checkpoint Size:每次 Checkpoint 的增量/全量大小- RocksDB 的
numBytesInMemTable,numBytesInBlockCache,numBytesInCompaction在 Metrics 中可查
Q4: State TTL 为什么不使用事件时间?
A: 可以使用,但有风险:
- 如果 Watermark 不推进(Idle Source / 巨大乱序),TTL 永远不会触发
- 状态会无限堆积,最终 OOM
- 建议:设置
cleanupIncrementally+cleanupInRocksdbCompactFilter兜底,即使 TTL 不触发也能通过访问时/Compaction 时清理
Q5: 为什么 RocksDB 状态下做了 KeyBy 后还需要序列化?
A: RocksDB 是 C++ 库(JNI 调用),数据必须从 JVM Heap 拷贝到 C++ 内存。这个过程:
- Java Object → TypeSerializer.serialize() → byte[]
- byte[] → JNI → C++ RocksDB → MemTable
- 读取时:C++ RocksDB → JNI → byte[] → TypeSerializer.deserialize() → Java Object
Q6: OperatorState 和 KeyedState 的区别?举一个具体的例子。
A: 以 Kafka Source 为例:
- OperatorState(Kafka Source):存储
(partition, offset)映射。因为 Kafka Partition 的分配不依赖 Key,而是依赖 Subtask。- Subtask 0 的 OperatorState: {partition-0: offset=1000, partition-1: offset=2000}
- Subtask 1 的 OperatorState: {partition-2: offset=1500, partition-3: offset=1800}
- KeyedState(Window 算子):存储每个 Key 的窗口聚合结果。
- Key "user_001" 的 KeyedState: {window[10:00]: count=5, window[10:01]: count=3}
- 这个状态按 Key 所属的 KeyGroup 分布在不同的 Subtask 中
下一步
- Window & Timer: doc-05-window-and-timer.md — Window 算子 + Timer 服务
- TM 内存模型: doc-10-taskmanager-memory.md — Managed Memory 与状态存储