聊聊flink的OperatorStateBackend

本文主要研究一下flink的OperatorStateBackend

OperatorStateBackend

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/OperatorStateBackend.java

/**
 * Interface that combines both, the user facing {@link OperatorStateStore} interface and the system interface
 * {@link Snapshotable}
 *
 */
public interface OperatorStateBackend extends
    OperatorStateStore,
    Snapshotable<SnapshotResult<OperatorStateHandle>, Collection<OperatorStateHandle>>,
    Closeable,
    Disposable {

    @Override
    void dispose();
}
  • OperatorStateBackend接口继承了OperatorStateStore、Snapshotable、Closeable、Disposable接口

OperatorStateStore

flink-core-1.7.0-sources.jar!/org/apache/flink/api/common/state/OperatorStateStore.java

/**
 * This interface contains methods for registering operator state with a managed store.
 */
@PublicEvolving
public interface OperatorStateStore {

    <K, V> BroadcastState<K, V> getBroadcastState(MapStateDescriptor<K, V> stateDescriptor) throws Exception;

    <S> ListState<S> getListState(ListStateDescriptor<S> stateDescriptor) throws Exception;

    <S> ListState<S> getUnionListState(ListStateDescriptor<S> stateDescriptor) throws Exception;

    Set<String> getRegisteredStateNames();

    Set<String> getRegisteredBroadcastStateNames();

    // -------------------------------------------------------------------------------------------
    //  Deprecated methods
    // -------------------------------------------------------------------------------------------

    @Deprecated
    <S> ListState<S> getOperatorState(ListStateDescriptor<S> stateDescriptor) throws Exception;


    @Deprecated
    <T extends Serializable> ListState<T> getSerializableListState(String stateName) throws Exception;
}
  • OperatorStateStore定义了getBroadcastState、getListState、getUnionListState方法用于create或restore BroadcastState或者ListState;同时也定义了getRegisteredStateNames、getRegisteredBroadcastStateNames用于返回当前注册的state的名称

Snapshotable

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/Snapshotable.java

/**
 * Interface for operators that can perform snapshots of their state.
 *
 * @param <S> Generic type of the state object that is created as handle to snapshots.
 * @param <R> Generic type of the state object that used in restore.
 */
@Internal
public interface Snapshotable<S extends StateObject, R> extends SnapshotStrategy<S> {

    /**
     * Restores state that was previously snapshotted from the provided parameters. Typically the parameters are state
     * handles from which the old state is read.
     *
     * @param state the old state to restore.
     */
    void restore(@Nullable R state) throws Exception;
}
  • Snapshotable接口继承了SnapshotStrategy接口,同时定义了restore方法用于restore state

SnapshotStrategy

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/SnapshotStrategy.java

/**
 * Interface for different snapshot approaches in state backends. Implementing classes should ideally be stateless or at
 * least threadsafe, i.e. this is a functional interface and is can be called in parallel by multiple checkpoints.
 *
 * @param <S> type of the returned state object that represents the result of the snapshot operation.
 */
@Internal
public interface SnapshotStrategy<S extends StateObject> {

    /**
     * Operation that writes a snapshot into a stream that is provided by the given {@link CheckpointStreamFactory} and
     * returns a @{@link RunnableFuture} that gives a state handle to the snapshot. It is up to the implementation if
     * the operation is performed synchronous or asynchronous. In the later case, the returned Runnable must be executed
     * first before obtaining the handle.
     *
     * @param checkpointId      The ID of the checkpoint.
     * @param timestamp         The timestamp of the checkpoint.
     * @param streamFactory     The factory that we can use for writing our state to streams.
     * @param checkpointOptions Options for how to perform this checkpoint.
     * @return A runnable future that will yield a {@link StateObject}.
     */
    @Nonnull
    RunnableFuture<S> snapshot(
        long checkpointId,
        long timestamp,
        @Nonnull CheckpointStreamFactory streamFactory,
        @Nonnull CheckpointOptions checkpointOptions) throws Exception;
}
  • SnapshotStrategy定义了snapshot方法,给不同的snapshot策略去实现,这里要求snapshot结果返回的类型是StateObject类型

AbstractSnapshotStrategy

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/AbstractSnapshotStrategy.java

/**
 * Abstract base class for implementing {@link SnapshotStrategy}, that gives a consistent logging across state backends.
 *
 * @param <T> type of the snapshot result.
 */
public abstract class AbstractSnapshotStrategy<T extends StateObject> implements SnapshotStrategy<SnapshotResult<T>> {

    private static final Logger LOG = LoggerFactory.getLogger(AbstractSnapshotStrategy.class);

    private static final String LOG_SYNC_COMPLETED_TEMPLATE = "{} ({}, synchronous part) in thread {} took {} ms.";
    private static final String LOG_ASYNC_COMPLETED_TEMPLATE = "{} ({}, asynchronous part) in thread {} took {} ms.";

    /** Descriptive name of the snapshot strategy that will appear in the log outputs and {@link #toString()}. */
    @Nonnull
    protected final String description;

    protected AbstractSnapshotStrategy(@Nonnull String description) {
        this.description = description;
    }

    /**
     * Logs the duration of the synchronous snapshot part from the given start time.
     */
    public void logSyncCompleted(@Nonnull Object checkpointOutDescription, long startTime) {
        logCompletedInternal(LOG_SYNC_COMPLETED_TEMPLATE, checkpointOutDescription, startTime);
    }

    /**
     * Logs the duration of the asynchronous snapshot part from the given start time.
     */
    public void logAsyncCompleted(@Nonnull Object checkpointOutDescription, long startTime) {
        logCompletedInternal(LOG_ASYNC_COMPLETED_TEMPLATE, checkpointOutDescription, startTime);
    }

    private void logCompletedInternal(
        @Nonnull String template,
        @Nonnull Object checkpointOutDescription,
        long startTime) {

        long duration = (System.currentTimeMillis() - startTime);

        LOG.debug(
            template,
            description,
            checkpointOutDescription,
            Thread.currentThread(),
            duration);
    }

    @Override
    public String toString() {
        return "SnapshotStrategy {" + description + "}";
    }
}
  • AbstractSnapshotStrategy是个抽象类,它没有实现SnapshotStrategy定义的snapshot方法,这里只是提供了logSyncCompleted方法打印debug信息

StateObject

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/StateObject.java

/**
 * Base of all handles that represent checkpointed state in some form. The object may hold
 * the (small) state directly, or contain a file path (state is in the file), or contain the
 * metadata to access the state stored in some external database.
 *
 * <p>State objects define how to {@link #discardState() discard state} and how to access the
 * {@link #getStateSize() size of the state}.
 * 
 * <p>State Objects are transported via RPC between <i>JobManager</i> and
 * <i>TaskManager</i> and must be {@link java.io.Serializable serializable} to support that.
 * 
 * <p>Some State Objects are stored in the checkpoint/savepoint metadata. For long-term
 * compatibility, they are not stored via {@link java.io.Serializable Java Serialization},
 * but through custom serializers.
 */
public interface StateObject extends Serializable {

    void discardState() throws Exception;

    long getStateSize();
}
  • StateObject继承了Serializable接口,因为会通过rpc在JobManager及TaskManager之间进行传输;这个接口定义了discardState及getStateSize方法,discardState用于清理资源,而getStateSize用于返回state的大小

StreamStateHandle

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/StreamStateHandle.java

/**
 * A {@link StateObject} that represents state that was written to a stream. The data can be read
 * back via {@link #openInputStream()}.
 */
public interface StreamStateHandle extends StateObject {

    /**
     * Returns an {@link FSDataInputStream} that can be used to read back the data that
     * was previously written to the stream.
     */
    FSDataInputStream openInputStream() throws IOException;
}
  • StreamStateHandle继承了StateObject接口,多定义了openInputStream方法

OperatorStateHandle

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/OperatorStateHandle.java

/**
 * Interface of a state handle for operator state.
 */
public interface OperatorStateHandle extends StreamStateHandle {

    /**
     * Returns a map of meta data for all contained states by their name.
     */
    Map<String, StateMetaInfo> getStateNameToPartitionOffsets();

    /**
     * Returns an input stream to read the operator state information.
     */
    @Override
    FSDataInputStream openInputStream() throws IOException;

    /**
     * Returns the underlying stream state handle that points to the state data.
     */
    StreamStateHandle getDelegateStateHandle();

    //......
}
  • OperatorStateHandle继承了StreamStateHandle,它多定义了getStateNameToPartitionOffsets、getDelegateStateHandle方法,其中getStateNameToPartitionOffsets提供了state name到可用partitions的offset的映射信息

OperatorStreamStateHandle

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/OperatorStreamStateHandle.java

/**
 * State handle for partitionable operator state. Besides being a {@link StreamStateHandle}, this also provides a
 * map that contains the offsets to the partitions of named states in the stream.
 */
public class OperatorStreamStateHandle implements OperatorStateHandle {

    private static final long serialVersionUID = 35876522969227335L;

    /**
     * unique state name -> offsets for available partitions in the handle stream
     */
    private final Map<String, StateMetaInfo> stateNameToPartitionOffsets;
    private final StreamStateHandle delegateStateHandle;

    public OperatorStreamStateHandle(
            Map<String, StateMetaInfo> stateNameToPartitionOffsets,
            StreamStateHandle delegateStateHandle) {

        this.delegateStateHandle = Preconditions.checkNotNull(delegateStateHandle);
        this.stateNameToPartitionOffsets = Preconditions.checkNotNull(stateNameToPartitionOffsets);
    }

    @Override
    public Map<String, StateMetaInfo> getStateNameToPartitionOffsets() {
        return stateNameToPartitionOffsets;
    }

    @Override
    public void discardState() throws Exception {
        delegateStateHandle.discardState();
    }

    @Override
    public long getStateSize() {
        return delegateStateHandle.getStateSize();
    }

    @Override
    public FSDataInputStream openInputStream() throws IOException {
        return delegateStateHandle.openInputStream();
    }

    @Override
    public StreamStateHandle getDelegateStateHandle() {
        return delegateStateHandle;
    }

    //......
}
  • OperatorStreamStateHandle实现了OperatorStateHandle接口,它定义了stateNameToPartitionOffsets属性(Map<String, StateMetaInfo>),而getStateNameToPartitionOffsets方法就是返回的stateNameToPartitionOffsets属性

SnapshotResult

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/SnapshotResult.java

/**
 * This class contains the combined results from the snapshot of a state backend:
 * <ul>
 *   <li>A state object representing the state that will be reported to the Job Manager to acknowledge the checkpoint.</li>
 *   <li>A state object that represents the state for the {@link TaskLocalStateStoreImpl}.</li>
 * </ul>
 *
 * Both state objects are optional and can be null, e.g. if there was no state to snapshot in the backend. A local
 * state object that is not null also requires a state to report to the job manager that is not null, because the
 * Job Manager always owns the ground truth about the checkpointed state.
 */
public class SnapshotResult<T extends StateObject> implements StateObject {

    private static final long serialVersionUID = 1L;

    /** An singleton instance to represent an empty snapshot result. */
    private static final SnapshotResult<?> EMPTY = new SnapshotResult<>(null, null);

    /** This is the state snapshot that will be reported to the Job Manager to acknowledge a checkpoint. */
    private final T jobManagerOwnedSnapshot;

    /** This is the state snapshot that will be reported to the Job Manager to acknowledge a checkpoint. */
    private final T taskLocalSnapshot;

    /**
     * Creates a {@link SnapshotResult} for the given jobManagerOwnedSnapshot and taskLocalSnapshot. If the
     * jobManagerOwnedSnapshot is null, taskLocalSnapshot must also be null.
     *
     * @param jobManagerOwnedSnapshot Snapshot for report to job manager. Can be null.
     * @param taskLocalSnapshot Snapshot for report to local state manager. This is optional and requires
     *                             jobManagerOwnedSnapshot to be not null if this is not also null.
     */
    private SnapshotResult(T jobManagerOwnedSnapshot, T taskLocalSnapshot) {

        if (jobManagerOwnedSnapshot == null && taskLocalSnapshot != null) {
            throw new IllegalStateException("Cannot report local state snapshot without corresponding remote state!");
        }

        this.jobManagerOwnedSnapshot = jobManagerOwnedSnapshot;
        this.taskLocalSnapshot = taskLocalSnapshot;
    }

    public T getJobManagerOwnedSnapshot() {
        return jobManagerOwnedSnapshot;
    }

    public T getTaskLocalSnapshot() {
        return taskLocalSnapshot;
    }

    @Override
    public void discardState() throws Exception {

        Exception aggregatedExceptions = null;

        if (jobManagerOwnedSnapshot != null) {
            try {
                jobManagerOwnedSnapshot.discardState();
            } catch (Exception remoteDiscardEx) {
                aggregatedExceptions = remoteDiscardEx;
            }
        }

        if (taskLocalSnapshot != null) {
            try {
                taskLocalSnapshot.discardState();
            } catch (Exception localDiscardEx) {
                aggregatedExceptions = ExceptionUtils.firstOrSuppressed(localDiscardEx, aggregatedExceptions);
            }
        }

        if (aggregatedExceptions != null) {
            throw aggregatedExceptions;
        }
    }

    @Override
    public long getStateSize() {
        return jobManagerOwnedSnapshot != null ? jobManagerOwnedSnapshot.getStateSize() : 0L;
    }

    @SuppressWarnings("unchecked")
    public static <T extends StateObject> SnapshotResult<T> empty() {
        return (SnapshotResult<T>) EMPTY;
    }

    public static <T extends StateObject> SnapshotResult<T> of(@Nullable T jobManagerState) {
        return jobManagerState != null ? new SnapshotResult<>(jobManagerState, null) : empty();
    }

    public static <T extends StateObject> SnapshotResult<T> withLocalState(
        @Nonnull T jobManagerState,
        @Nonnull T localState) {
        return new SnapshotResult<>(jobManagerState, localState);
    }
}
  • SnapshotResult类实现了StateObject接口,它包装了snapshot的结果,这里包括jobManagerOwnedSnapshot、taskLocalSnapshot;它实现的discardState方法,调用了jobManagerOwnedSnapshot及taskLocalSnapshot的discardState方法;getStateSize方法则返回的是jobManagerOwnedSnapshot的stateSize

DefaultOperatorStateBackend

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java

/**
 * Default implementation of OperatorStateStore that provides the ability to make snapshots.
 */
@Internal
public class DefaultOperatorStateBackend implements OperatorStateBackend {
    
    private static final Logger LOG = LoggerFactory.getLogger(DefaultOperatorStateBackend.class);

    /**
     * The default namespace for state in cases where no state name is provided
     */
    public static final String DEFAULT_OPERATOR_STATE_NAME = "_default_";

    /**
     * Map for all registered operator states. Maps state name -> state
     */
    private final Map<String, PartitionableListState<?>> registeredOperatorStates;

    /**
     * Map for all registered operator broadcast states. Maps state name -> state
     */
    private final Map<String, BackendWritableBroadcastState<?, ?>> registeredBroadcastStates;

    /**
     * CloseableRegistry to participate in the tasks lifecycle.
     */
    private final CloseableRegistry closeStreamOnCancelRegistry;

    /**
     * Default serializer. Only used for the default operator state.
     */
    private final JavaSerializer<Serializable> javaSerializer;

    /**
     * The user code classloader.
     */
    private final ClassLoader userClassloader;

    /**
     * The execution configuration.
     */
    private final ExecutionConfig executionConfig;

    /**
     * Flag to de/activate asynchronous snapshots.
     */
    private final boolean asynchronousSnapshots;

    /**
     * Map of state names to their corresponding restored state meta info.
     *
     * <p>TODO this map can be removed when eager-state registration is in place.
     * TODO we currently need this cached to check state migration strategies when new serializers are registered.
     */
    private final Map<String, StateMetaInfoSnapshot> restoredOperatorStateMetaInfos;

    /**
     * Map of state names to their corresponding restored broadcast state meta info.
     */
    private final Map<String, StateMetaInfoSnapshot> restoredBroadcastStateMetaInfos;

    /**
     * Cache of already accessed states.
     *
     * <p>In contrast to {@link #registeredOperatorStates} and {@link #restoredOperatorStateMetaInfos} which may be repopulated
     * with restored state, this map is always empty at the beginning.
     *
     * <p>TODO this map should be moved to a base class once we have proper hierarchy for the operator state backends.
     *
     * @see <a href="https://issues.apache.org/jira/browse/FLINK-6849">FLINK-6849</a>
     */
    private final HashMap<String, PartitionableListState<?>> accessedStatesByName;

    private final Map<String, BackendWritableBroadcastState<?, ?>> accessedBroadcastStatesByName;

    private final AbstractSnapshotStrategy<OperatorStateHandle> snapshotStrategy;

    public DefaultOperatorStateBackend(
        ClassLoader userClassLoader,
        ExecutionConfig executionConfig,
        boolean asynchronousSnapshots) {

        this.closeStreamOnCancelRegistry = new CloseableRegistry();
        this.userClassloader = Preconditions.checkNotNull(userClassLoader);
        this.executionConfig = executionConfig;
        this.javaSerializer = new JavaSerializer<>();
        this.registeredOperatorStates = new HashMap<>();
        this.registeredBroadcastStates = new HashMap<>();
        this.asynchronousSnapshots = asynchronousSnapshots;
        this.accessedStatesByName = new HashMap<>();
        this.accessedBroadcastStatesByName = new HashMap<>();
        this.restoredOperatorStateMetaInfos = new HashMap<>();
        this.restoredBroadcastStateMetaInfos = new HashMap<>();
        this.snapshotStrategy = new DefaultOperatorStateBackendSnapshotStrategy();
    }

    @Override
    public Set<String> getRegisteredStateNames() {
        return registeredOperatorStates.keySet();
    }

    @Override
    public Set<String> getRegisteredBroadcastStateNames() {
        return registeredBroadcastStates.keySet();
    }

    @Override
    public void close() throws IOException {
        closeStreamOnCancelRegistry.close();
    }

    @Override
    public void dispose() {
        IOUtils.closeQuietly(closeStreamOnCancelRegistry);
        registeredOperatorStates.clear();
        registeredBroadcastStates.clear();
    }

    // -------------------------------------------------------------------------------------------
    //  State access methods
    // -------------------------------------------------------------------------------------------

    @SuppressWarnings("unchecked")
    @Override
    public <K, V> BroadcastState<K, V> getBroadcastState(final MapStateDescriptor<K, V> stateDescriptor) throws StateMigrationException {
        //......
    }

    @Override
    public <S> ListState<S> getListState(ListStateDescriptor<S> stateDescriptor) throws Exception {
        return getListState(stateDescriptor, OperatorStateHandle.Mode.SPLIT_DISTRIBUTE);
    }

    @Override
    public <S> ListState<S> getUnionListState(ListStateDescriptor<S> stateDescriptor) throws Exception {
        return getListState(stateDescriptor, OperatorStateHandle.Mode.UNION);
    }

    @Nonnull
    @Override
    public RunnableFuture<SnapshotResult<OperatorStateHandle>> snapshot(
        long checkpointId,
        long timestamp,
        @Nonnull CheckpointStreamFactory streamFactory,
        @Nonnull CheckpointOptions checkpointOptions) throws Exception {

        long syncStartTime = System.currentTimeMillis();

        RunnableFuture<SnapshotResult<OperatorStateHandle>> snapshotRunner =
            snapshotStrategy.snapshot(checkpointId, timestamp, streamFactory, checkpointOptions);

        snapshotStrategy.logSyncCompleted(streamFactory, syncStartTime);
        return snapshotRunner;
    }

    //......
}
  • DefaultOperatorStateBackend实现了OperatorStateBackend接口
  • getRegisteredStateNames方法返回的是registeredOperatorStates.keySet();getRegisteredBroadcastStateNames方法返回的是registeredBroadcastStates.keySet(),可以看到这两个都是基于内存的Map来实现的
  • close方法主要是调用closeStreamOnCancelRegistry的close方法;dispose方法也会关闭closeStreamOnCancelRegistry,同时清空registeredOperatorStates及registeredBroadcastStates
  • getListState及getUnionListState方法都调用了getListState(ListStateDescriptor<S> stateDescriptor,OperatorStateHandle.Mode mode)方法
  • snapshot方法使用的snapshotStrategy是DefaultOperatorStateBackendSnapshotStrategy

DefaultOperatorStateBackend.getListState

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java

    private <S> ListState<S> getListState(
            ListStateDescriptor<S> stateDescriptor,
            OperatorStateHandle.Mode mode) throws StateMigrationException {

        Preconditions.checkNotNull(stateDescriptor);
        String name = Preconditions.checkNotNull(stateDescriptor.getName());

        @SuppressWarnings("unchecked")
        PartitionableListState<S> previous = (PartitionableListState<S>) accessedStatesByName.get(name);
        if (previous != null) {
            checkStateNameAndMode(
                    previous.getStateMetaInfo().getName(),
                    name,
                    previous.getStateMetaInfo().getAssignmentMode(),
                    mode);
            return previous;
        }

        // end up here if its the first time access after execution for the
        // provided state name; check compatibility of restored state, if any
        // TODO with eager registration in place, these checks should be moved to restore()

        stateDescriptor.initializeSerializerUnlessSet(getExecutionConfig());
        TypeSerializer<S> partitionStateSerializer = Preconditions.checkNotNull(stateDescriptor.getElementSerializer());

        @SuppressWarnings("unchecked")
        PartitionableListState<S> partitionableListState = (PartitionableListState<S>) registeredOperatorStates.get(name);

        if (null == partitionableListState) {
            // no restored state for the state name; simply create new state holder

            partitionableListState = new PartitionableListState<>(
                new RegisteredOperatorStateBackendMetaInfo<>(
                    name,
                    partitionStateSerializer,
                    mode));

            registeredOperatorStates.put(name, partitionableListState);
        } else {
            // has restored state; check compatibility of new state access

            checkStateNameAndMode(
                    partitionableListState.getStateMetaInfo().getName(),
                    name,
                    partitionableListState.getStateMetaInfo().getAssignmentMode(),
                    mode);

            StateMetaInfoSnapshot restoredSnapshot = restoredOperatorStateMetaInfos.get(name);
            RegisteredOperatorStateBackendMetaInfo<S> metaInfo =
                new RegisteredOperatorStateBackendMetaInfo<>(restoredSnapshot);

            // check compatibility to determine if state migration is required
            TypeSerializer<S> newPartitionStateSerializer = partitionStateSerializer.duplicate();

            @SuppressWarnings("unchecked")
            TypeSerializerSnapshot<S> stateSerializerSnapshot = Preconditions.checkNotNull(
                (TypeSerializerSnapshot<S>) restoredSnapshot.getTypeSerializerConfigSnapshot(StateMetaInfoSnapshot.CommonSerializerKeys.VALUE_SERIALIZER));

            TypeSerializerSchemaCompatibility<S> stateCompatibility =
                stateSerializerSnapshot.resolveSchemaCompatibility(newPartitionStateSerializer);
            if (stateCompatibility.isIncompatible()) {
                throw new StateMigrationException("The new state serializer for operator state must not be incompatible.");
            }

            partitionableListState.setStateMetaInfo(
                new RegisteredOperatorStateBackendMetaInfo<>(name, newPartitionStateSerializer, mode));
        }

        accessedStatesByName.put(name, partitionableListState);
        return partitionableListState;
    }
  • 从registeredOperatorStates获取对应PartitionableListState,没有的话则创建,有的话则检查下兼容性,然后往partitionableListState设置stateMetaInfo

DefaultOperatorStateBackendSnapshotStrategy

flink-runtime_2.11-1.7.0-sources.jar!/org/apache/flink/runtime/state/DefaultOperatorStateBackend.java

    /**
     * Snapshot strategy for this backend.
     */
    private class DefaultOperatorStateBackendSnapshotStrategy extends AbstractSnapshotStrategy<OperatorStateHandle> {

        protected DefaultOperatorStateBackendSnapshotStrategy() {
            super("DefaultOperatorStateBackend snapshot");
        }

        @Nonnull
        @Override
        public RunnableFuture<SnapshotResult<OperatorStateHandle>> snapshot(
            final long checkpointId,
            final long timestamp,
            @Nonnull final CheckpointStreamFactory streamFactory,
            @Nonnull final CheckpointOptions checkpointOptions) throws IOException {

            if (registeredOperatorStates.isEmpty() && registeredBroadcastStates.isEmpty()) {
                return DoneFuture.of(SnapshotResult.empty());
            }

            final Map<String, PartitionableListState<?>> registeredOperatorStatesDeepCopies =
                new HashMap<>(registeredOperatorStates.size());
            final Map<String, BackendWritableBroadcastState<?, ?>> registeredBroadcastStatesDeepCopies =
                new HashMap<>(registeredBroadcastStates.size());

            ClassLoader snapshotClassLoader = Thread.currentThread().getContextClassLoader();
            Thread.currentThread().setContextClassLoader(userClassloader);
            try {
                // eagerly create deep copies of the list and the broadcast states (if any)
                // in the synchronous phase, so that we can use them in the async writing.

                if (!registeredOperatorStates.isEmpty()) {
                    for (Map.Entry<String, PartitionableListState<?>> entry : registeredOperatorStates.entrySet()) {
                        PartitionableListState<?> listState = entry.getValue();
                        if (null != listState) {
                            listState = listState.deepCopy();
                        }
                        registeredOperatorStatesDeepCopies.put(entry.getKey(), listState);
                    }
                }

                if (!registeredBroadcastStates.isEmpty()) {
                    for (Map.Entry<String, BackendWritableBroadcastState<?, ?>> entry : registeredBroadcastStates.entrySet()) {
                        BackendWritableBroadcastState<?, ?> broadcastState = entry.getValue();
                        if (null != broadcastState) {
                            broadcastState = broadcastState.deepCopy();
                        }
                        registeredBroadcastStatesDeepCopies.put(entry.getKey(), broadcastState);
                    }
                }
            } finally {
                Thread.currentThread().setContextClassLoader(snapshotClassLoader);
            }

            AsyncSnapshotCallable<SnapshotResult<OperatorStateHandle>> snapshotCallable =
                new AsyncSnapshotCallable<SnapshotResult<OperatorStateHandle>>() {

                    @Override
                    protected SnapshotResult<OperatorStateHandle> callInternal() throws Exception {

                        CheckpointStreamFactory.CheckpointStateOutputStream localOut =
                            streamFactory.createCheckpointStateOutputStream(CheckpointedStateScope.EXCLUSIVE);
                        registerCloseableForCancellation(localOut);

                        // get the registered operator state infos ...
                        List<StateMetaInfoSnapshot> operatorMetaInfoSnapshots =
                            new ArrayList<>(registeredOperatorStatesDeepCopies.size());

                        for (Map.Entry<String, PartitionableListState<?>> entry :
                            registeredOperatorStatesDeepCopies.entrySet()) {
                            operatorMetaInfoSnapshots.add(entry.getValue().getStateMetaInfo().snapshot());
                        }

                        // ... get the registered broadcast operator state infos ...
                        List<StateMetaInfoSnapshot> broadcastMetaInfoSnapshots =
                            new ArrayList<>(registeredBroadcastStatesDeepCopies.size());

                        for (Map.Entry<String, BackendWritableBroadcastState<?, ?>> entry :
                            registeredBroadcastStatesDeepCopies.entrySet()) {
                            broadcastMetaInfoSnapshots.add(entry.getValue().getStateMetaInfo().snapshot());
                        }

                        // ... write them all in the checkpoint stream ...
                        DataOutputView dov = new DataOutputViewStreamWrapper(localOut);

                        OperatorBackendSerializationProxy backendSerializationProxy =
                            new OperatorBackendSerializationProxy(operatorMetaInfoSnapshots, broadcastMetaInfoSnapshots);

                        backendSerializationProxy.write(dov);

                        // ... and then go for the states ...

                        // we put BOTH normal and broadcast state metadata here
                        int initialMapCapacity =
                            registeredOperatorStatesDeepCopies.size() + registeredBroadcastStatesDeepCopies.size();
                        final Map<String, OperatorStateHandle.StateMetaInfo> writtenStatesMetaData =
                            new HashMap<>(initialMapCapacity);

                        for (Map.Entry<String, PartitionableListState<?>> entry :
                            registeredOperatorStatesDeepCopies.entrySet()) {

                            PartitionableListState<?> value = entry.getValue();
                            long[] partitionOffsets = value.write(localOut);
                            OperatorStateHandle.Mode mode = value.getStateMetaInfo().getAssignmentMode();
                            writtenStatesMetaData.put(
                                entry.getKey(),
                                new OperatorStateHandle.StateMetaInfo(partitionOffsets, mode));
                        }

                        // ... and the broadcast states themselves ...
                        for (Map.Entry<String, BackendWritableBroadcastState<?, ?>> entry :
                            registeredBroadcastStatesDeepCopies.entrySet()) {

                            BackendWritableBroadcastState<?, ?> value = entry.getValue();
                            long[] partitionOffsets = {value.write(localOut)};
                            OperatorStateHandle.Mode mode = value.getStateMetaInfo().getAssignmentMode();
                            writtenStatesMetaData.put(
                                entry.getKey(),
                                new OperatorStateHandle.StateMetaInfo(partitionOffsets, mode));
                        }

                        // ... and, finally, create the state handle.
                        OperatorStateHandle retValue = null;

                        if (unregisterCloseableFromCancellation(localOut)) {

                            StreamStateHandle stateHandle = localOut.closeAndGetHandle();

                            if (stateHandle != null) {
                                retValue = new OperatorStreamStateHandle(writtenStatesMetaData, stateHandle);
                            }

                            return SnapshotResult.of(retValue);
                        } else {
                            throw new IOException("Stream was already unregistered.");
                        }
                    }

                    @Override
                    protected void cleanupProvidedResources() {
                        // nothing to do
                    }

                    @Override
                    protected void logAsyncSnapshotComplete(long startTime) {
                        if (asynchronousSnapshots) {
                            logAsyncCompleted(streamFactory, startTime);
                        }
                    }
                };

            final FutureTask<SnapshotResult<OperatorStateHandle>> task =
                snapshotCallable.toAsyncSnapshotFutureTask(closeStreamOnCancelRegistry);

            if (!asynchronousSnapshots) {
                task.run();
            }

            return task;
        }
    }
  • DefaultOperatorStateBackendSnapshotStrategy继承了AbstractSnapshotStrategy,它实现的snapshot方法主要是创建registeredOperatorStatesDeepCopies及registeredBroadcastStatesDeepCopies,然后通过AsyncSnapshotCallable来实现
  • AsyncSnapshotCallable抽象类实现了Callable接口的call方法,该方法会调用callInternal方法,然后再执行logAsyncSnapshotComplete方法
  • AsyncSnapshotCallable的callInternal方法返回的是SnapshotResult<OperatorStateHandle>,它里头主要是将registeredOperatorStatesDeepCopies及registeredBroadcastStatesDeepCopies的数据写入到CheckpointStreamFactory(比如MemCheckpointStreamFactory).CheckpointStateOutputStream及writtenStatesMetaData,最后通过CheckpointStateOutputStream的closeAndGetHandle返回的stateHandle及writtenStatesMetaData创建OperatorStreamStateHandle返回

小结

  • OperatorStateBackend接口继承了OperatorStateStore、Snapshotable、Closeable、Disposable接口
  • OperatorStateStore定义了getBroadcastState、getListState、getUnionListState方法用于create或restore BroadcastState或者ListState;同时也定义了getRegisteredStateNames、getRegisteredBroadcastStateNames用于返回当前注册的state的名称;DefaultOperatorStateBackend实现了OperatorStateStore接口,getRegisteredStateNames方法返回的是registeredOperatorStates.keySet();getRegisteredBroadcastStateNames方法返回的是registeredBroadcastStates.keySet()(registeredOperatorStates及registeredBroadcastStates这两个都是内存的Map);getListState及getUnionListState方法都调用了getListState(ListStateDescriptor<S> stateDescriptor,OperatorStateHandle.Mode mode)方法
  • Snapshotable接口继承了SnapshotStrategy接口,同时定义了restore方法用于restore state;SnapshotStrategy定义了snapshot方法,给不同的snapshot策略去实现,这里要求snapshot结果返回的类型是StateObject类型;AbstractSnapshotStrategy是个抽象类,它没有实现SnapshotStrategy定义的snapshot方法,这里只是提供了logSyncCompleted方法打印debug信息
  • DefaultOperatorStateBackend实现了Snapshotable接口,snapshot方法使用的snapshotStrategy是DefaultOperatorStateBackendSnapshotStrategy;DefaultOperatorStateBackendSnapshotStrategy继承了AbstractSnapshotStrategy,它实现的snapshot方法主要是创建registeredOperatorStatesDeepCopies及registeredBroadcastStatesDeepCopies,然后通过AsyncSnapshotCallable来实现,它里头主要是将registeredOperatorStatesDeepCopies及registeredBroadcastStatesDeepCopies的数据写入到CheckpointStreamFactory(比如MemCheckpointStreamFactory).CheckpointStateOutputStream及writtenStatesMetaData
  • Snapshotable接口要求source的泛型为StateObject类型,StateObject继承了Serializable接口,因为会通过rpc在JobManager及TaskManager之间进行传输;OperatorStateBackend继承Snapshotable接口时,指定source为SnapshotResult<OperatorStateHandle>,而result的为Collection<OperatorStateHandle>类型
  • StreamStateHandle继承了StateObject接口,多定义了openInputStream方法;OperatorStateHandle继承了StreamStateHandle,它多定义了getStateNameToPartitionOffsets、getDelegateStateHandle方法,其中getStateNameToPartitionOffsets提供了state name到可用partitions的offset的映射信息;OperatorStreamStateHandle实现了OperatorStateHandle接口,它定义了stateNameToPartitionOffsets属性(Map<String,StateMetaInfo>),而getStateNameToPartitionOffsets方法就是返回的stateNameToPartitionOffsets属性
  • SnapshotResult类实现了StateObject接口,它包装了snapshot的结果,这里包括jobManagerOwnedSnapshot、taskLocalSnapshot;它实现的discardState方法,调用了jobManagerOwnedSnapshot及taskLocalSnapshot的discardState方法;getStateSize方法则返回的是jobManagerOwnedSnapshot的stateSize

doc

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 213,417评论 6 492
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 90,921评论 3 387
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 158,850评论 0 349
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 56,945评论 1 285
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 66,069评论 6 385
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 50,188评论 1 291
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,239评论 3 412
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 37,994评论 0 268
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 44,409评论 1 304
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 36,735评论 2 327
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 38,898评论 1 341
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 34,578评论 4 336
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,205评论 3 317
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 30,916评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,156评论 1 267
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 46,722评论 2 363
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 43,781评论 2 351

推荐阅读更多精彩内容

  • Spring Cloud为开发人员提供了快速构建分布式系统中一些常见模式的工具(例如配置管理,服务发现,断路器,智...
    卡卡罗2017阅读 134,644评论 18 139
  • 序 本文主要研究一下flink的CheckpointScheduler CheckpointCoordinator...
    go4it阅读 1,853评论 0 0
  • 序 本文主要研究一下flink的SourceFunction 实例 这里通过addSource方法来添加自定义的S...
    go4it阅读 8,703评论 0 3
  • 序 本文主要研究一下flink的CheckpointedFunction 实例 这个BufferingSink实现...
    go4it阅读 5,205评论 0 1
  • 随着大数据 2.0 时代悄然到来,大数据从简单的批处理扩展到了实时处理、流处理、交互式查询和机器学习应用。近年来涌...
    Java大生阅读 2,106评论 0 6