Amazon Managed Service for Apache Flink is the AWS service that runs your Apache Flink streaming job without you operating a Flink cluster. It was launched as Kinesis Data Analytics for Apache Flink and renamed in August 2023, which is why the API, the CLI command and the IAM actions still say kinesisanalyticsv2 and its capacity unit is still called a Kinesis Processing Unit. The older sibling, Kinesis Data Analytics for SQL applications, is gone: AWS stopped those applications on 15 October 2025 and began deleting them on 27 January 2026, pointing users at Managed Flink or Managed Flink Studio instead.

This article is about running the managed service, not Flink internals: what AWS runs for you, how KPUs turn into parallelism and cost, how autoscaling decides, how snapshots protect state across deploys, the Flink 2.2 upgrade, and the failure modes that page people.

What AWS runs for you

An application in Managed Flink is three things you supply and one thing AWS supplies. You supply the code, a JAR for Java or a ZIP for Python, stored in S3; a service execution role that the job assumes to reach its sources and sinks; and configuration, which includes parallelism, checkpointing, logging and named groups of runtime properties. AWS supplies a dedicated Flink cluster: a JobManager that schedules work and coordinates checkpoints, and TaskManager capacity sliced into KPUs that run your operators.

You never see the hosts. You cannot SSH in, install agents or edit flink-conf.yaml beyond the settings the service exposes. In exchange the service patches the runtime, replaces failed hosts, stores snapshots durably and restarts the job from the last checkpoint when a task fails. If you attach the application to a VPC, the service places elastic network interfaces in your subnets so the job can reach private resources such as an MSK cluster; it then needs a NAT gateway or VPC endpoints to reach public AWS APIs like Kinesis or CloudWatch.

What you own and what AWS runs in Managed Service for Apache FlinkKinesis / MSKsourcesS3 code bucketJAR or ZIP, versionedRuntime propertiesPropertyGroupsService-managed cluster (per application)JobManagerorchestration KPU (billed)KPU 11 vCPU, 4 GB, 50 GBKPU NParallelism / PerKPUCheckpoints (every 60 s default)Snapshots = managed savepointsrecordsdeployconfigSinksS3, Kinesis, DynamoDBCloudWatchmetrics + logsYour VPCENIs if configuredresultsControl plane: kinesisanalyticsv2 APICreate, Update, Start, Stop, Snapshot, RollbackYou own the code, the IAM role, the parallelism settings and the state schema. AWS owns hosts, Flink processes and patching.
The managed boundary: everything inside the grey box is AWS-operated; everything outside it is configuration you own.

KPUs, parallelism and the maxParallelism trap

A KPU is the unit of both capacity and billing: one vCPU, 4 GB of memory and 50 GB of running application storage, which is the local disk RocksDB uses for operator state. Two settings decide how many you get. Parallelism is the default number of parallel instances of every operator, default 1. ParallelismPerKPU is how many of those parallel tasks share one KPU, default 1 and maximum 8. The service allocates Parallelism / ParallelismPerKPU KPUs for task work and charges one more KPU for orchestration.

The ratio is the interesting knob. A CPU-bound job wants one task per vCPU, so leave ParallelismPerKPU at 1. A job that mostly waits on asynchronous calls can pack several tasks per KPU, but packing also splits the 4 GB and 50 GB between them, so state-heavy operators suffer first.

SettingDefaultLimitWhat it changes
Parallelism1ParallelismPerKPU x KPU quotaOperator parallelism, and with the ratio, KPU count
ParallelismPerKPU18Tasks per KPU; memory and disk per task
KPUs per application-64 (quota, raisable)Hard ceiling for scaling
maxParallelism128 if parallelism is 128 or lessSet in codeNumber of key groups; ceiling for rescaling with state

The last row is the trap. Flink hashes keys into a fixed number of key groups, the job's maximum parallelism. If you never set it and start with parallelism of 128 or less, every operator gets 128. Autoscaling will not scale past it, a manual update past it fails to start, and changing it later means you cannot restore from earlier snapshots. Set env.setMaxParallelism(...) deliberately on day one, to a value comfortably above your largest plausible parallelism.

Worked example: sizing a clickstream job

Suppose a clickstream arrives on a Kinesis data stream, and the job enriches each event with a user segment from DynamoDB, counts events per user in five-minute tumbling windows and writes results to S3. Load testing in staging, with the same record mix as production, shows that one task keeps up with roughly 1,500 records per second at 60 percent CPU, and that most of its time is spent in the asynchronous DynamoDB lookup. Production peaks at 15,000 records per second.

Ten tasks would run at about 60 percent at peak, with headroom for catch-up after a restart. Because the job is I/O-bound, you try ParallelismPerKPU of 2, which gives Parallelism 10 and 5 task KPUs plus 1 orchestration KPU, 6 KPUs billed per hour. A few gigabytes of window state fits easily in the 250 GB of running storage. You set maxParallelism to 720, which divides evenly by many parallelism values, and you re-run the load test at the packed ratio, because the CPU number only transfers if the tasks really are waiting most of the time.

Per-task throughput depends entirely on your code and records, so measure it rather than borrowing a records-per-KPU figure. The job skeleton on the Flink 1.20 runtime with connector 5.x:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setMaxParallelism(720);                       // fixed for the life of the state

Properties io = KinesisAnalyticsRuntime.getApplicationProperties().get("io");  // runtime properties
String inputArn = io.getProperty("input.stream.arn");
String outputUri = io.getProperty("output.s3.uri");

Configuration srcCfg = new Configuration();
srcCfg.set(KinesisSourceConfigOptions.STREAM_INITIAL_POSITION,
           KinesisSourceConfigOptions.InitialPosition.TRIM_HORIZON);

KinesisStreamsSource<String> source = KinesisStreamsSource.<String>builder()
    .setStreamArn(inputArn)
    .setSourceConfig(srcCfg)
    .setDeserializationSchema(new SimpleStringSchema())
    .build();

DataStream<Click> clicks = env
    .fromSource(source, WatermarkStrategy.<String>forBoundedOutOfOrderness(Duration.ofSeconds(20)),
                "clicks-source")
    .uid("clicks-source")                         // stable IDs map state across deploys
    .map(Click::parse).uid("parse");

DataStream<Click> enriched = AsyncDataStream.unorderedWait(
        clicks, new SegmentLookup(), 500, TimeUnit.MILLISECONDS, 100).uid("segment-lookup");

enriched.keyBy(Click::userId)
    .window(TumblingEventTimeWindows.of(Duration.ofMinutes(5)))
    .aggregate(new CountPerUser()).uid("count-5m")
    .map(Object::toString)
    .sinkTo(FileSink.forRowFormat(new Path(outputUri),
                                  new SimpleStringEncoder<String>("UTF-8")).build())
    .uid("s3-sink");

env.execute("clickstream-enrichment");

Deploys, snapshots and rollback

Every change to a running application goes through UpdateApplication, which takes the current application version ID as an optimistic lock, so two pipelines cannot silently overwrite each other. When snapshots are enabled, an update stops the job with a snapshot, applies the new code or configuration and restarts from that snapshot. Starting a stopped application takes a restore type: SKIP_RESTORE_FROM_SNAPSHOT for a clean start, RESTORE_FROM_LATEST_SNAPSHOT for normal operation, or RESTORE_FROM_CUSTOM_SNAPSHOT with a snapshot name to roll back to a known point. A deployment script that takes a named snapshot first gives you that known point:

APP=clickstream-enrichment
TAG=pre-deploy-$(date +%Y%m%d%H%M)

aws kinesisanalyticsv2 create-application-snapshot --application-name $APP --snapshot-name $TAG

VER=$(aws kinesisanalyticsv2 describe-application --application-name $APP \
      --query 'ApplicationDetail.ApplicationVersionId' --output text)

aws kinesisanalyticsv2 update-application --application-name $APP \
  --current-application-version-id $VER \
  --application-configuration-update '{"ApplicationCodeConfigurationUpdate":
    {"CodeContentUpdate":{"S3ContentLocationUpdate":{"FileKeyUpdate":"jobs/clicks-1.42.jar"}}}}'

# Roll back by name if the new version misbehaves:
# aws kinesisanalyticsv2 stop-application --application-name $APP
# aws kinesisanalyticsv2 start-application --application-name $APP --run-configuration \
#   '{"ApplicationRestoreConfiguration":{"ApplicationRestoreType":"RESTORE_FROM_CUSTOM_SNAPSHOT",
#     "SnapshotName":"'$TAG'"}}'

Leave AllowNonRestoredState false so an accidental uid change fails loudly instead of silently dropping state, and consider enabling system rollback, which returns to the previous version when an update fails to start. Neither catches a version that starts fine and computes the wrong answer; the named snapshot does.

Checkpoints versus snapshots

Checkpoints and snapshots are different tools. Checkpoints are Flink's automatic, frequent, incremental copies of state used to recover from a task failure; the service owns them and you cannot restore from an arbitrary one. With the DEFAULT checkpoint configuration the service forces checkpointing on, every 60,000 ms, with at least 5,000 ms between checkpoints, and it ignores whatever your code sets. To change the interval you must set the configuration type to CUSTOM in the application configuration, not in code.

Snapshots are the service's name for savepoints: durable, named, taken on stop, update and scaling, or on demand through CreateApplicationSnapshot. They are what makes redeploys stateful. Their contract with your code is the operator uid and the state serializer. Change a uid and that operator's state no longer maps; change a POJO's fields carelessly and the serializer may refuse the old bytes. For the mechanics of barriers and alignment, see Flink checkpoint design; for choosing between heap and RocksDB state, see Flink state backends compared; and for savepoint discipline in general, Flink savepoints.

How automatic scaling decides

Automatic scaling is on by default and it is simpler than people assume. The service watches the maximum containerCPUUtilization in one-minute datapoints. Fifteen consecutive datapoints at or above 75 percent trigger a scale-up that doubles current parallelism, and with it the KPU count. Three hundred and sixty consecutive datapoints, six hours, below 10 percent trigger a scale-down that halves parallelism, rounded up. It never goes below the Parallelism you configured and never above the job's maximum parallelism or KPU quota.

Consequences follow. First, scaling is a restart: the application enters the AUTOSCALING status and comes back at the new parallelism, and AWS documents downtime while it does. Second, the signal is CPU only. A job that falls behind because a sink is throttling or because it waits on a slow lookup can run at low CPU while lag grows, and the autoscaler will not react; it may even scale down. Third, doubling is coarse. Going from 10 to 20 tasks when 12 would do doubles the bill until six quiet hours pass.

AWS publishes samples for scheduled scaling and for scaling on other metrics such as Kinesis millisBehindLatest; both still restart the job, but at a better moment. If restarts are unacceptable, disable autoscaling and size for peak.

Metrics that answer real questions

Managed Flink publishes metrics to CloudWatch at a level you choose, from application down to operator, task or parallelism level; finer levels give more insight at higher CloudWatch cost. A small set answers most questions:

QuestionMetricAlarm on
Are we keeping up?millisBehindLatest (Kinesis) or consumer lag (Kafka)Sustained growth, not a single spike
Is the job healthy?numRestarts (Flink 2.2), fullRestarts on 1.xAny increase outside a deploy
Is state safe?lastCheckpointDuration, numberOfFailedCheckpointsDuration near the interval; any failures
Is it sized right?containerCPUUtilization, containerMemoryUtilizationCPU high with lag growing; memory near full
Where is the bottleneck?backPressuredTimeMsPerSecond per operatorThe first busy operator upstream of back-pressure

A restart loop shows up in the CloudWatch log stream as the same stack trace every checkpoint interval.

Moving from Flink 1.x to the 2.2 runtime

AWS added Flink 2.2 in March 2026, the service's first major-version runtime, and supports in-place upgrades that keep configuration, logs, metrics and, when compatible, state. Older runtimes are being retired: 1.6, 1.8 and 1.11 were stopped in July 2025 and 1.13 support ended on 16 October 2025. Plan the 1.x to 2.2 move as a migration, because several changes break applications that ran fine before:

  • APIs removed. The DataSet API, the Scala API and the legacy SourceFunction/SinkFunction interfaces are gone. Kinesis jobs move to KinesisStreamsSource and KinesisStreamsSink from connector 6.0.0-2.0.
  • Runtimes. Java 17 and Python 3.12 are the defaults; Java 11 and Python 3.8 are removed.
  • State compatibility. Kryo moves from 2.24 to 5.6 and POJOs holding collections may not restore; Avro and Protobuf state is unaffected. Read the state compatibility guide and test a restore from a production snapshot copy before upgrading.
  • Sandboxing. The root filesystem is read-only except /tmp, and IMDS calls other than the credential and region paths are blocked, so code that asks for its instance ID will fail.
  • Gaps. Studio notebooks do not run on 2.2, and at the time of writing JDBC, OpenSearch and Prometheus connectors had no 2.x release.

Failure modes

These are the incidents that recur, with the signature to look for and the fix.

  • Poison record restart loop. One malformed record throws in a map function; the job restarts from the last checkpoint, reads the same record and throws again. numRestarts climbs while output stops. Catch parse errors and route bad records to a side output written to S3 rather than letting them fail the task.
  • Checkpoints timing out under back-pressure. Barriers queue behind slow records, checkpoint duration climbs toward the timeout, and eventually the job restarts and replays more data than before. Fix the slow operator, often a synchronous external call, before raising the interval.
  • Running storage full. RocksDB state outgrows the 50 GB per KPU, usually because a keyed state has no TTL. Add state TTL or windowing so keys expire; adding KPUs only buys time.
  • Enhanced fan-out known issues. AWS documents that KinesisStreamsSource with enhanced fan-out in connectors 5.0.0 and 6.0.0 may fail when a stream is resharded (FLINK-37648), and, combined with KinesisStreamsSink under back-pressure, can deadlock so that a force stop and start are needed (FLINK-34071). If you depend on EFO, rehearse a reshard in staging and alarm on output stalls.

Trade-offs

Managed Flink competes with three alternatives. Self-managed Flink on EKS with the Flink Kubernetes operator gives full control of configuration, metric reporters and versions, at the cost of running Kubernetes and Flink upgrades yourself. Lambda with a Kinesis or MSK trigger is cheaper and simpler for stateless per-record work, but has no managed keyed state, event-time windows or exactly-once aggregation across records. Managed Flink sits between: real Flink semantics with very little operations.

Its costs are a floor of two KPUs per application even when idle, restarts on every scale or deploy, a CPU-only autoscaler, and a restricted configuration surface that tightens on 2.2. Many small applications each pay an orchestration KPU. Before choosing, confirm that your source really needs stateful streaming: the sources themselves are covered in Amazon Kinesis and Amazon MSK.

What to do next

  1. Set maxParallelism explicitly in code, and give every operator a stable uid.
  2. Load-test one task with production-shaped records and derive Parallelism and ParallelismPerKPU from the measurement.
  3. Decide whether CPU-driven autoscaling fits; if lag is your real signal, disable it and use scheduled or custom scaling.
  4. Add a named snapshot step to your deployment script and practise a rollback with RESTORE_FROM_CUSTOM_SNAPSHOT.
  5. Alarm on lag growth, restarts, failed checkpoints and checkpoint duration, not just CPU.
  6. Route unparseable records to a side output instead of letting them fail the job.
  7. Add state TTL to every keyed state that can grow without bound.
  8. If you are on a 1.x runtime, inventory DataSet, Scala, legacy connector and Kryo-serialised POJO use before planning the 2.2 upgrade.
Key takeaway: Managed Flink removes cluster operations but not Flink discipline. Size from a measured per-task rate, fix maxParallelism and operator uids on day one, treat every scale and deploy as a snapshot-and-restart, alarm on lag rather than CPU, and plan the move to Flink 2.2 as a migration with a tested state restore.