Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@
import org.slf4j.LoggerFactory;

import javax.annotation.Nonnull;
import javax.annotation.Nullable;

import java.io.IOException;
import java.time.Duration;
Expand Down Expand Up @@ -149,15 +150,60 @@ protected void log(InfoLogLevel infoLogLevel, String logMsg) {

public void configCompactFilter(
@Nonnull StateDescriptor<?, ?> stateDesc, TypeSerializer<?> stateSerializer) {
StateTtlConfig ttlConfig = stateDesc.getTtlConfig();
// A V1 list state is registered with a ListSerializer<TtlValue<E>>; the compaction filter
// works on the elements, so unwrap it here.
TypeSerializer<?> listElementSerializer =
stateDesc instanceof ListStateDescriptor
? ((ListSerializer<?>) stateSerializer).getElementSerializer()
: null;
configCompactFilter(
stateDesc.getName(),
stateDesc.getTtlConfig(),
listElementSerializer,
stateDesc instanceof MapStateDescriptor);
}

/**
* Configures the TTL compaction filter for a state created through the State V2 API. This is
* the counterpart of {@link #configCompactFilter(StateDescriptor, TypeSerializer)}: registering
* the filter factory on the column family alone (see {@link
* #setAndRegisterCompactFilterIfStateTtlV2}) is not sufficient, the native filter stays
* disabled until it receives its {@link FlinkCompactionFilter.Config}.
*/
public void configCompactFilter(
@Nonnull org.apache.flink.api.common.state.v2.StateDescriptor<?> stateDesc,
TypeSerializer<?> stateSerializer) {
org.apache.flink.api.common.state.v2.StateDescriptor.Type type = stateDesc.getType();
// Unlike V1, a V2 ListStateDescriptor carries the element serializer itself (for a TTL
// state: TtlSerializer<E>), and ForStListState stores the elements delimited with that
// serializer. So the registered state serializer is already the list element serializer.
TypeSerializer<?> listElementSerializer =
type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.LIST
? stateSerializer
: null;
configCompactFilter(
stateDesc.getStateId(),
stateDesc.getTtlConfig(),
listElementSerializer,
type == org.apache.flink.api.common.state.v2.StateDescriptor.Type.MAP);
}

/**
* @param listElementSerializer the element serializer of a list state (the TTL-wrapped element
* serializer), or {@code null} if the state is not a list state.
*/
private void configCompactFilter(
String stateName,
StateTtlConfig ttlConfig,
@Nullable TypeSerializer<?> listElementSerializer,
boolean isMapState) {
if (ttlConfig.isEnabled() && ttlConfig.getCleanupStrategies().inRocksdbCompactFilter()) {
FlinkCompactionFilterFactory compactionFilterFactory =
compactionFilterFactories.get(stateDesc.getName());
compactionFilterFactories.get(stateName);
Preconditions.checkNotNull(compactionFilterFactory);
long ttl = ttlConfig.getTimeToLive().toMillis();

ColumnFamilyOptions columnFamilyOptions =
columnFamilyOptionsMap.get(stateDesc.getName());
ColumnFamilyOptions columnFamilyOptions = columnFamilyOptionsMap.get(stateName);
Preconditions.checkNotNull(columnFamilyOptions);

StateTtlConfig.RocksdbCompactFilterCleanupStrategy rocksdbCompactFilterCleanupStrategy =
Expand All @@ -178,9 +224,8 @@ public void configCompactFilter(
}

FlinkCompactionFilter.Config config;
if (stateDesc instanceof ListStateDescriptor) {
TypeSerializer<?> elemSerializer =
((ListSerializer<?>) stateSerializer).getElementSerializer();
if (listElementSerializer != null) {
TypeSerializer<?> elemSerializer = listElementSerializer;
int len = elemSerializer.getLength();
if (len > 0) {
config =
Expand All @@ -195,7 +240,7 @@ public void configCompactFilter(
queryTimeAfterNumEntries,
new ListElementFilterFactory<>(elemSerializer.duplicate()));
}
} else if (stateDesc instanceof MapStateDescriptor) {
} else if (isMapState) {
config = FlinkCompactionFilter.Config.createForMap(ttl, queryTimeAfterNumEntries);
} else {
config = FlinkCompactionFilter.Config.createForValue(ttl, queryTimeAfterNumEntries);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@
import org.forstdb.ColumnFamilyHandle;
import org.forstdb.ColumnFamilyOptions;
import org.forstdb.RocksDB;
import org.forstdb.RocksDBException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -291,6 +292,14 @@ public <N, S extends InternalKeyedState, SV> S createStateInternal(
Tuple2<ColumnFamilyHandle, RegisteredKeyValueStateBackendMetaInfo<N, SV>> registerResult =
tryRegisterKvStateInformation(stateDesc, namespaceSerializer);

// The compaction filter factory is attached to the column family when the state is
// registered, but the native filter only becomes active once it receives the TTL
// configuration. Mirror what the sync backend does after registration.
if (ttlCompactFiltersManager != null) {
ttlCompactFiltersManager.configCompactFilter(
stateDesc, registerResult.f1.getStateSerializer());
}

ColumnFamilyHandle columnFamilyHandle = registerResult.f0;

switch (stateDesc.getType()) {
Expand Down Expand Up @@ -624,6 +633,13 @@ public boolean isSafeToReuseKVState() {
return true;
}

@VisibleForTesting
public void compactState(StateDescriptor<?> stateDesc) throws RocksDBException {
ForStOperationUtils.ForStKvStateInfo kvStateInfo =
kvStateInformation.get(stateDesc.getStateId());
db.compactRange(kvStateInfo.columnFamilyHandle);
}

@VisibleForTesting
Path getLocalBasePath() {
return optionsContainer.getPathContainer().getLocalBasePath();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,17 @@ public static <K> ForStKeyedStateBackend<K> createKeyedStateBackend(
TypeSerializer<K> keySerializer,
Collection<KeyedStateHandle> stateHandles)
throws IOException {
return createKeyedStateBackend(
forStStateBackend, env, keySerializer, stateHandles, TtlTimeProvider.DEFAULT);
}

public static <K> ForStKeyedStateBackend<K> createKeyedStateBackend(
ForStStateBackend forStStateBackend,
Environment env,
TypeSerializer<K> keySerializer,
Collection<KeyedStateHandle> stateHandles,
TtlTimeProvider ttlTimeProvider)
throws IOException {

return forStStateBackend.createAsyncKeyedStateBackend(
new KeyedStateBackendParametersImpl<>(
Expand All @@ -51,7 +62,7 @@ public static <K> ForStKeyedStateBackend<K> createKeyedStateBackend(
1,
new KeyGroupRange(0, 0),
env.getTaskKvStateRegistry(),
TtlTimeProvider.DEFAULT,
ttlTimeProvider,
new UnregisteredMetricsGroup(),
(name, value) -> {},
stateHandles,
Expand Down
Loading