All skills
aws avatar

/developing-applications-on-managed-service-for-apache-flink

@222ce56

MANDATORY for Flink or Amazon Managed Service for Apache Flink (MSF) questions. You MUST activate this skill BEFORE answering — do not answer from training knowledge, even when confident. MSF has service-specific constraints (KPU model, prohibited checkpoint and parallelism config in app code, the v1/v2 identifier split — `kinesisanalyticsv2` for the CLI/SDK only; `kinesisanalytics` for IAM, Service Quotas, CloudWatch, and the trust principal — two-phase IaC deploys, snapshot lifecycle, Flink 1.x→2.x migration) that override generic Flink knowledge.

Use this Skill: https://skilld.dev/gh/aws/agent-toolkit-for-aws/developing-applications-on-managed-service-for-apache-flink

This session only. Nothing lands on disk.

referencesstate-management.md

≈1.2k tokens on demand. Your agent reads this file only when SKILL.md points to it.

State Management Best Practices

Overview

This guide covers state management best practices for Managed Service for Apache Flink applications, including efficient state usage, TTL configuration, state type selection, and Managed Service for Apache Flink-specific state management considerations.

For general development patterns and application structure, see best-practices.md. For serialization guidance, see serialization-guide.md.

Code examples in this guide use Flink 2.2 APIs by default, which are also compatible with Flink 1.20 unless noted otherwise. See flink-2x-migration.md for the complete migration reference.

Efficient State Usage with Managed Service for Apache Flink

  • Estimate state size ahead of time to ensure state will remain bounded over time.
  • Enable state TTL to ensure state gets cleaned up automatically if not cleaned up manually.
  • Use the correct state type for each use case, and perform updates to state in performant way (e.g. make updates to a map key, rather than replacing the entire map).
  • AVOID: Storing large objects or unbounded collections in state

Pick the Right State Type — MapState vs ValueState<Map>

Use MapState<K, V> whenever you need per-key updates inside a logical map. Storing a Map<K, V> inside ValueState<Map<K, V>> and reading-mutating-writing it on every event is O(map size) per access — RocksDB has to deserialize every entry, your code mutates one, and the whole map gets re-serialized and written back. MapState is O(1) per put/get/remove: each map entry maps to its own RocksDB key, so only the touched entry is read or written.

There is also a state-migration consequence specific to MSF: nested generic collections inside ValueState (e.g. ValueState<Map<String, MyType>>) typically fall back to Kryo serialization, and Kryo-serialized state does not migrate from Flink 1.x to 2.x. MapState uses a dedicated MapSerializer that does carry across the upgrade. See serialization-guide.md for the wider Kryo guidance and flink-2x-migration.md for the migration impact.

public class OptimizedKeyedProcessor extends KeyedProcessFunction<String, Event, Result> {
    
    // Use appropriate state types
    private transient ValueState<EventAggregate> aggregateState;
    
    @Override
    public void open(OpenContext openContext) throws Exception {
        ValueStateDescriptor<EventAggregate> aggregateDescriptor = 
            new ValueStateDescriptor<>("aggregate", EventAggregate.class);
        
        // TTL configuration - application-level concern
        aggregateDescriptor.enableTimeToLive(StateTtlConfig.newBuilder(Duration.ofHours(24))
            .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
            .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
            .build());
        
        aggregateState = getRuntimeContext().getState(aggregateDescriptor);
    }
    
    @Override
    public void processElement(Event event, Context ctx, Collector<Result> out) throws Exception {
        // Efficient state access patterns
        EventAggregate current = aggregateState.value();
        if (current == null) {
            current = new EventAggregate();
        }
        
        // Update state efficiently
        current.update(event);
        aggregateState.update(current);
        
        // Use timers for handling events that occur after the input event - e.g. in this case we want to trigger the output an hour after the input occurs
        ctx.timerService().registerEventTimeTimer(event.getTimestamp() + 3600000); // 1 hour
    }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<Result> out) throws Exception {
        // Handle timer firing - cleanup or emit final results
        EventAggregate current = aggregateState.value();
        if (current != null) {
            // Emit final result or perform cleanup
            out.collect(new Result(ctx.getCurrentKey(), current.getFinalValue()));
            // Clear state after processing to re-initialize if needed
            aggregateState.clear();
        }
    }
}

Managed Service for Apache Flink State Management

Managed Service for Apache Flink service handles:

  • State backend configuration (RocksDB with S3 for checkpoints/savepoints)
  • Checkpoint storage and retention
  • Savepoint management through console
  • State size monitoring and alerting Application code should NOT configure these aspects

Source: SKILL.md on GitHub

No alerts2mo3 checks · Risk SAFE
  • Gen Agent Trust Hub2mo

    This skill provides comprehensive guidance for developing applications on Amazon Managed Service for Apache Flink (MSF). The content is focused on technical documentation, best practices, and code examples aligned with AWS and Apache Flink standards. No security issues were detected during the analysis.

  • Socket2mo

    No alerts

  • Snyk2mo

    Risk: LOW · No issues

Signed by skilld at 222ce56. This ties the file your Agent reads to that commit on GitHub. It does not review the instructions.

Last checked against GitHub yesterday.

Activeupdated 2 months ago
version
2
Other metadata
Triggers — activate on any of
Flink, MSF, Managed Flink, KinesisAnalytics(V2), KPU, ParallelismPerKPU, savepoint, checkpoint, operator UID, FlinkKinesisConsumer, KinesisStreamsSource, KafkaSource, IcebergSink, EFO, CreateApplication, UpdateApplication, CreateApplicationSnapshot, Kryo, RocksDB, Iceberg streaming, EXACTLY_ONCE, watermark, CDC binlog/WAL, Glue/S3 Tables, AWS/KinesisAnalytics CloudWatch.

README badge

README badge for aws/agent-toolkit-for-aws/developing-applications-on-managed-service-for-apache-flink