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.

referencesserialization-guide.md

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

Serialization Best Practices

Overview

This guide covers serialization best practices for Managed Service for Apache Flink applications, including the performance hierarchy of serializer types, POJO and Tuple usage, Avro and Protobuf integration, Kryo avoidance, state serialization considerations, and common anti-patterns.

For general development patterns and application structure, see best-practices.md. For state management guidance, see state-management.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.

Performance Hierarchy: Choose the Right Serializer

Flink serialization performance (fastest to slowest):

  1. Flink Tuples/Rows - Fastest (direct field access, no reflection)
  2. POJOs - Fast (~30% slower than tuples, supports schema evolution)
  3. Protobuf - Good performance (~30% slower than POJOs)
  4. Avro Specific - Moderate (~50% slower than POJOs)
  5. Avro Generic/Thrift - Slower (~70% slower than POJOs)
  6. Kryo - Avoid (50%+ performance penalty). Kryo registration convenience methods are removed from StreamExecutionEnvironment in Flink 2.x; avoid Kryo entirely.

POJO Serialization for Managed Service for Apache Flink

// Recommended: Flink POJO for optimal performance with schema evolution
public class OptimizedEvent {
    // All fields must be public or have public getters/setters
    public String eventId;
    public long timestamp;
    public String userId;
    public EventType type;
    
    // Required: public no-argument constructor
    public OptimizedEvent() {}
    
    public OptimizedEvent(String eventId, long timestamp, String userId, EventType type) {
        this.eventId = eventId;
        this.timestamp = timestamp;
        this.userId = userId;
        this.type = type;
    }
}

// Enum types work well with POJO serialization
public enum EventType {
    USER_ACTION, SYSTEM_EVENT, ERROR_EVENT
}

Tuple Types for Maximum Performance

// Use when performance is critical and schema evolution is not needed
DataStream<Tuple4<String, Long, String, Integer>> events = source
    .map(event -> Tuple4.of(event.getId(), event.getTimestamp(), 
                           event.getUserId(), event.getCount()));

// Access fields by position (f0, f1, f2, f3)
events.keyBy(tuple -> tuple.f2) // Key by userId (f2)
      .process(new TupleProcessor());

Avro for External Integration

// Use Avro when integrating with external systems or when advanced schema evolution is needed
public class AvroEventProcessor extends ProcessFunction<SpecificRecordBase, ProcessedEvent> {
    
    @Override
    public void processElement(SpecificRecordBase avroEvent, Context ctx, 
                              Collector<ProcessedEvent> out) {
        if (avroEvent instanceof UserEvent) {
            UserEvent userEvent = (UserEvent) avroEvent;
            ProcessedEvent result = new ProcessedEvent();
            result.setUserId(userEvent.getUserId().toString());
            result.setTimestamp(userEvent.getTimestamp());
            out.collect(result);
        }
    }
}

// Configure Avro serialization
// Note: enableForceAvro() is available in Flink 1.20 but removed in 2.x.
// For Flink 2.2, use AvroTypeInfo explicitly in state descriptors instead.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Flink 1.20 only:
// env.getConfig().enableForceAvro();

Protobuf Integration

The recommended approach for Protobuf is to convert Protobuf messages to Flink POJOs at the ingestion boundary. This avoids Kryo entirely, giving you fast serialization, schema evolution support, and state compatibility across Flink major versions.

// Recommended: Convert Protobuf to POJO at the boundary, avoiding Kryo entirely
DataStream<MyEvent> events = protobufSource
    .map(proto -> new MyEvent(
        proto.getEventId(),
        proto.getTimestamp(),
        proto.getUserId()));
// MyEvent is a Flink POJO (public fields + no-arg constructor) — fast serialization, schema evolution, no Kryo

Legacy note — not recommended for new applications: If you must use Protobuf objects directly in state, you can register them with Kryo via env.getConfig(). However, Kryo has a 50%+ performance penalty and Kryo-serialized state does not migrate from Flink 1.x to 2.x. Convenience registration methods on StreamExecutionEnvironment are removed in Flink 2.x.

// Not recommended — use POJO conversion instead:
env.getConfig().registerTypeWithKryoSerializer(
    MyProtobufMessage.class,
    ProtobufSerializer.class
);

Avoiding Kryo Fallbacks

Why Kryo Fallback Matters on MSF

If Flink can't recognize a type as a POJO/Tuple/Avro/Protobuf, it silently falls back to Kryo. On MSF this has three consequences worth treating as blockers, not warnings:

  • ~50% performance penalty vs. POJO serialization, plus larger serialized objects on the wire and in state. On a high-throughput keyed pipeline this dominates per-record cost.
  • Larger checkpoint and shuffle bytes. Inflated checkpoint size lengthens the checkpoint window and pushes more data across cross-AZ network paths inside MSF.
  • Kryo-serialized state does not migrate from Flink 1.x to 2.x. This is a hard blocker for in-place version upgrades — see flink-2x-migration.md for the migration path. Plan to eliminate Kryo before the 1→2 upgrade, not after.

Fail Fast in Development

// Monitor for Kryo fallbacks in logs - these indicate performance issues
// Log message: "Class ... cannot be used as a POJO type because not all fields are valid POJO fields"

// To detect Kryo usage, disable it temporarily during development
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.getConfig().disableGenericTypes(); // Throws exception if Kryo would be used
// This will fail with: "Generic types have been disabled in the ExecutionConfig"

Run with disableGenericTypes() enabled locally as part of every PR build so Kryo fallbacks fail the build, not production.

Last Resort: Kryo Type Registration

Warning: Prefer converting to Flink POJOs or Tuples instead of registering Kryo serializers. Kryo-serialized state does not migrate across Flink major versions. Convenience registration methods on StreamExecutionEnvironment are removed in Flink 2.x — use env.getConfig() methods instead. Use this only when migrating away from Kryo is not yet feasible.

// Last resort — register frequently used types to avoid class name serialization overhead IF you use Kryo
// In Flink 2.x, use env.getConfig() methods (env-level convenience methods are removed)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

env.getConfig().registerKryoType(CustomEvent.class);
env.getConfig().registerKryoType(ProcessingResult.class);

env.getConfig().registerTypeWithKryoSerializer(
    ComplexObject.class, 
    CustomKryoSerializer.class
);

State Serialization Considerations

// For state objects, prioritize schema evolution support
public class StatefulProcessor extends KeyedProcessFunction<String, Event, Result> {
    
    // Use POJO or Avro for state that needs to evolve
    private transient ValueState<EventAggregate> aggregateState; // POJO - good performance + evolution
    
    // Use primitive types for simple state
    private transient ValueState<Long> counterState; // Primitive - fastest
    
    @Override
    public void open(OpenContext openContext) throws Exception {
        // POJO state descriptor - supports schema evolution
        ValueStateDescriptor<EventAggregate> aggregateDescriptor = 
            new ValueStateDescriptor<>("aggregate", EventAggregate.class);
        aggregateState = getRuntimeContext().getState(aggregateDescriptor);
        
        // Primitive state descriptor - fastest serialization
        ValueStateDescriptor<Long> counterDescriptor = 
            new ValueStateDescriptor<>("counter", Long.class);
        counterState = getRuntimeContext().getState(counterDescriptor);
    }
}

Anti-Patterns: Serialization Performance Killers

// AVOID: Default Java Serialization (extremely slow)
public class SlowEvent implements Serializable {
    // Java serialization is 10x+ slower than POJO serialization
}

// AVOID: Complex nested objects without proper POJO structure
public class BadEvent {
    private Map<String, Object> data; // Generic Object causes Kryo fallback
    private List<SomeInterface> items; // Interface types cause Kryo fallback
}

// AVOID: Missing no-argument constructor (causes Kryo fallback)
public class InvalidPOJO {
    public String field;
    
    // Missing no-arg constructor - will use Kryo instead of POJO serializer
    public InvalidPOJO(String field) {
        this.field = field;
    }
}

// AVOID: Private fields without getters/setters (causes Kryo fallback)
public class AlmostPOJO {
    private String secretField; // No getter/setter - not a valid POJO field
    public String publicField;   // This is fine
    
    public AlmostPOJO() {} // Constructor is correct
}

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