如何为Apache Flink开发原始键控状态?自定义布隆过滤器状态问询
Got it, I totally get why you're looking for a more efficient way here—storing the entire BloomFilter instance in ValueState means you have to read/write the whole object every time you update it, which gets pretty costly especially with large filters or remote state backends. The raw state API in Flink is indeed under-documented, so let's break this down step by step, including both a custom raw state implementation and a simpler alternative that might fit your needs.
Why ValueState<BloomFilter> is inefficient
Every time you add an element to the bloom filter, you have to:
- Fetch the entire
BloomFilterinstance from state - Modify it in memory
- Serialize the whole instance and write it back to state
This full round-trip is wasteful for bloom filters, since we only need to flip individual bits rather than touch the entire object. Raw state lets us work directly with byte arrays, cutting out the overhead of serializing/deserializing the entire BloomFilter instance.
Option 1: Implement Custom Raw State
Raw state is Flink's lowest-level state API, giving you full control over how state is stored and manipulated. Here's how to build a custom raw state for your bloom filter:
Step 1: Define the Raw State Class
This class handles the bloom filter logic directly on a byte buffer, avoiding object serialization overhead:
public class BloomFilterRawState implements State { private final ByteBuffer bitmap; private final int numHashFunctions; private final int bitSize; public BloomFilterRawState(int bitSize, int numHashFunctions) { this.bitSize = bitSize; this.numHashFunctions = numHashFunctions; this.bitmap = ByteBuffer.allocate(bitSize / 8); } // Add element by flipping relevant bits directly in the byte buffer public void add(byte[] element) { for (int i = 0; i < numHashFunctions; i++) { int hash = Math.abs(MurmurHash3.hash(element, i)) % bitSize; int byteIndex = hash / 8; int bitIndex = hash % 8; bitmap.put(byteIndex, (byte) (bitmap.get(byteIndex) | (1 << bitIndex))); } } // Check if element might exist public boolean mightContain(byte[] element) { for (int i = 0; i < numHashFunctions; i++) { int hash = Math.abs(MurmurHash3.hash(element, i)) % bitSize; int byteIndex = hash / 8; int bitIndex = hash % 8; if ((bitmap.get(byteIndex) & (1 << bitIndex)) == 0) { return false; } } return true; } @Override public void clear() { bitmap.clear(); Arrays.fill(bitmap.array(), (byte) 0); } // Serialize state to byte array for checkpointing public byte[] serialize() { return bitmap.array(); } // Deserialize state from byte array (for recovery) public static BloomFilterRawState deserialize(byte[] data, int bitSize, int numHashFunctions) { BloomFilterRawState state = new BloomFilterRawState(bitSize, numHashFunctions); state.bitmap.put(data); state.bitmap.flip(); return state; } }
Step 2: Create a State Descriptor
This tells Flink how to serialize/deserialize your custom state:
public class BloomFilterStateDescriptor extends StateDescriptor<BloomFilterRawState, byte[]> { private final int bitSize; private final int numHashFunctions; public BloomFilterStateDescriptor(String name, int bitSize, int numHashFunctions) { super(name, TypeInformation.of(byte[].class), null); this.bitSize = bitSize; this.numHashFunctions = numHashFunctions; } @Override public StateSerializer<BloomFilterRawState> createSerializer(ExecutionConfig executionConfig) { return new StateSerializer<BloomFilterRawState>() { @Override public byte[] serialize(BloomFilterRawState state) throws IOException { return state.serialize(); } @Override public BloomFilterRawState deserialize(byte[] data) throws IOException { return BloomFilterRawState.deserialize(data, bitSize, numHashFunctions); } @Override public BloomFilterRawState copy(BloomFilterRawState from) throws IOException { BloomFilterRawState copy = new BloomFilterRawState(bitSize, numHashFunctions); copy.bitmap.put(from.bitmap.array()); copy.bitmap.flip(); return copy; } }; } }
Step 3: Use the Custom State in an Operator
Now you can use this state in your process function:
public class BloomFilterProcessFunction extends RichProcessFunction<String, String> { private BloomFilterRawState bloomFilterState; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // Configure bloom filter: 1MB bitmap, 5 hash functions BloomFilterStateDescriptor descriptor = new BloomFilterStateDescriptor( "custom-bloom-filter", 1024 * 1024 * 8, 5 ); bloomFilterState = getRuntimeContext().getState(descriptor); } @Override public void processElement(String value, Context ctx, Collector<String> out) throws Exception { byte[] elementBytes = value.getBytes(StandardCharsets.UTF_8); if (!bloomFilterState.mightContain(elementBytes)) { bloomFilterState.add(elementBytes); out.collect(value); // Output only new elements } } }
Option 2: Simpler Alternative - ValueState<byte[]>
If you don't want to dive into raw state, you can store just the bloom filter's bitmap as a byte[] in ValueState. This avoids serializing the entire BloomFilter object, which is already a big win:
// Scala implementation example class BloomFilterFunction extends RichProcessFunction[String, String] { private var bitmapState: ValueState[Array[Byte]] = _ private val bitSize = 1024 * 1024 * 8 // 1MB bitmap private val numHashFunctions = 5 override def open(parameters: Configuration): Unit = { val descriptor = new ValueStateDescriptor[Array[Byte]]( "bloom-filter-bitmap", createTypeInformation[Array[Byte]] ) bitmapState = getRuntimeContext.getState(descriptor) // Initialize bitmap if state is empty if (bitmapState.value() == null) { bitmapState.update(new Array[Byte](bitSize / 8)) } } override def processElement(value: String, ctx: Context, out: Collector[String]): Unit = { val bytes = value.getBytes(StandardCharsets.UTF_8) var elementExists = true val bitmap = bitmapState.value() for (i <- 0 until numHashFunctions) { val hash = Math.abs(MurmurHash3.hash(bytes, i)) % bitSize val byteIdx = hash / 8 val bitIdx = hash % 8 if ((bitmap(byteIdx) & (1 << bitIdx)) == 0) { elementExists = false // Flip the bit to 1 bitmap(byteIdx) = (bitmap(byteIdx) | (1 << bitIdx)).toByte } } if (!elementExists) { bitmapState.update(bitmap) out.collect(value) } } }
Key Notes
- Compatibility: When using raw state, you're responsible for handling schema changes (e.g., if you change the bit size or hash function count, old checkpoints won't be compatible).
- Dynamic Scaling: If you need a bloom filter that can grow, raw state will require more complex logic to handle bitmap resizing while maintaining consistency.
- Memory: For very large bloom filters, consider using Flink's managed memory instead of heap memory to avoid GC pressure—this requires working with off-heap byte buffers.
内容的提问来源于stack exchange,提问作者Moein Hosseini

