Use Cases of Flink CheckpointedFunction

Viewed 516

While going through the Flink official documentation, I came across CheckpointedFunction. Wondering why and when would you use this function. I am currently working on a stateful Flink job that heavily relies on ProcessFunction to save state in RocksDB. Just wondering if CheckpointedFunction is better than the ProcessFunction.

2 Answers

CheckpointedFunction is for cases where you need to work with state that should be managed by Flink and included in checkpoints, but where you aren't working with a KeyedStream and so you cannot use keyed state like you would in a KeyedProcessFunction.

The most common use cases of CheckpointedFunction are in sources and sinks.

In addition to the answer of @David I have another use case in which I don't use CheckpointedFunction with the source or sink. I do use it in a ProcessFunction where I want to count (programmatically) how many times my job has restarted. I use MyProcessFunction and CheckpointedFunction and I update ListState<Long> restarts when the job restarts. I use this state on the integration tests to ensure that the job was restarted upon a failure. I based my example on the Flink checkpoint example for Sinks.

public class MyProcessFunction<V> extends ProcessFunction<V, V> implements CheckpointedFunction {
    ...

    private transient ListState<Long> restarts;

    @Override
    public void snapshotState(FunctionSnapshotContext context) throws Exception { ... }

    @Override
    public void initializeState(FunctionInitializationContext context) throws Exception {
        restarts = context.getOperatorStateStore().getListState(new ListStateDescriptor<Long>("restarts", Long.class));

        if (context.isRestored()) {
            List<Long> restoreList = Lists.newArrayList(restarts.get());
            if (restoreList == null || restoreList.isEmpty()) {
                restarts.add(1L);
                System.out.println("restarts: 1");
            } else {
                Long max = Collections.max(restoreList);
                System.out.println("restarts: " + max);
                restarts.add(max + 1);
            }
        } else {
            System.out.println("restarts: never restored");
        }
    }

    @Override
    public void open(Configuration parameters) throws Exception { ... }

    @Override
    public void processElement(V value, Context ctx, Collector<V> out) throws Exception { ... }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<V> out) throws Exception { ... }
}
Related