Apache Flink is a distributed engine for computations over streams of events that need to remember things: running totals, the last known state of a device, whether a customer has placed six orders in a minute. It keeps that memory as local state next to the computation, orders events by the time they happened rather than the time they arrived, and periodically snapshots everything so a crash costs a replay instead of a wrong answer.

This page is the on-ramp. It explains the four ideas you need, how to choose between Flink's APIs, and walks one realistic job written both in SQL and in Java. It also covers how to deploy the job and keep it upgradable, and what changed in Flink 2.x. The runtime internals (slots, chaining, graph translation, barrier alignment) have their own page, Apache Flink architecture in depth, and this one links to the deeper pages where they apply.

Advertisement

The model in one picture

A Flink job: sources, keyed operators with local state, sinks, and periodic snapshotsKafka: orderspartitions 0..11source + watermarksevent time from o.tskeyBy(customerId)hash to key groupsinksalerts, spend tablekeyed process / window operatorsper-key state + event-time timersrecords for key kresultslocal stateheap or RocksDBtimer servicefires on watermarkJobManagerschedules, coordinates checkpointsTaskManagersrun operator instances in slotscheckpoint storageS3 / HDFS / GCSdeploysnapshots
A job is a dataflow graph. keyBy routes every record for a key to one operator instance, which owns that key's state and timers; the JobManager snapshots all state to durable storage on a schedule.

A Flink program describes a dataflow graph. Sources read from Kafka, Kinesis, files or change-data-capture streams; operators transform records; sinks write results. When you submit the program, the JobManager turns the graph into parallel tasks and deploys them onto TaskManagers, which are JVM processes that each offer a number of slots. An operator with parallelism 12 runs as 12 instances, usually spread across machines.

The step that makes Flink different from a stateless stream mapper is keyBy. It hashes a key, such as a customer id, and routes every record with that key to the same operator instance. That instance can then keep state for the key in a local store, either on the JVM heap or in an embedded RocksDB database on local disk.

Four ideas you need before writing code

Streams are unbounded, and bounded data is a special case. The same DataStream program can read a Kafka topic forever or a set of Parquet files that ends. Flink 2.x removed the old DataSet API, so batch work also runs through DataStream or SQL with bounded sources.

State is part of the computation. A counter, a map of open sessions or a buffered window is state. Flink owns it: it scopes it to the current key, stores it in the configured backend and includes it in snapshots. You never write it to an external database to survive a restart. Flink state, in depth covers the state types and what each access costs.

Time comes from the event. With event time, each record carries a timestamp, and the source emits watermarks. A watermark is a claim: no record older than time t is still expected. Windows and timers fire when the watermark passes their end, so the result does not depend on how fast or in what order the data arrived. Bounded out-of-orderness of five seconds means the watermark trails the largest timestamp seen by five seconds; events later than that are late and need explicit handling.

Snapshots make failure cheap. Every checkpoint interval, the JobManager injects markers into the sources. Each operator snapshots its state when the marker passes, and the checkpoint completes when every task has reported in. Together with the source offsets stored in the same snapshot, this gives a consistent cut. After a crash, Flink restores every operator from the last completed checkpoint and rewinds the sources to the matching offsets, so each event affects state exactly once. Effects on external systems are exactly once only if the sink is transactional or idempotent; exactly-once semantics in streaming explains that boundary.

Advertisement

Choosing an API

LayerYou writeGood forYou give up
Flink SQL / Table APISQL with window table functions, joins, MATCH_RECOGNIZEAggregations, joins, enrichment, ETL, anything an analyst can readFine control over state layout and timers; state compatibility across query changes is limited
DataStream APIJava (or Python) operators: map, keyBy, window, connectCustom logic, custom serialization, connectors with side effectsOptimizer; you choose types and state yourself
KeyedProcessFunctionPer-event callback with keyed state and timersRules engines, timeouts, custom windows, state machinesEverything is manual, including cleanup

Start with SQL. A surprising share of streaming work is aggregation and joining, and the SQL planner picks efficient operators, handles retractions and lets non-Java teams own the logic. Drop to DataStream when you need behaviour SQL cannot express, and to a process function when the logic is fundamentally a per-key state machine with timeouts.

One caveat decides many real choices. Changing a SQL query usually changes the operator graph the planner generates, and the old state may not map onto it. With DataStream you control operator identity and state names, so evolving a long-lived stateful job is more predictable. If a job holds weeks of state that is expensive to rebuild, weigh that before committing to SQL.

Worked example, part 1: spend per minute in SQL

A payments team has an orders topic with about 20,000 events per second and two needs: a per-customer spend total for each minute, written to a table that dashboards read, and an alert when a customer places more than five orders within a minute. The first is a textbook aggregation, so it is SQL.

CREATE TABLE orders (
  order_id    STRING,
  customer_id STRING,
  amount      DECIMAL(10, 2),
  ts          TIMESTAMP(3),
  WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'orders',
  'properties.bootstrap.servers' = 'kafka:9092',
  'properties.group.id' = 'spend-per-minute',
  'scan.startup.mode' = 'earliest-offset',
  'format' = 'json'
);

INSERT INTO spend_per_minute
SELECT window_start, window_end, customer_id,
       SUM(amount) AS spend, COUNT(*) AS order_count
FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
GROUP BY window_start, window_end, customer_id;

The WATERMARK clause declares event time and tolerates five seconds of disorder. The TUMBLE table function assigns each row to a one-minute window, and because the aggregation groups by the window bounds, the planner emits each window once, when the watermark passes its end, as an append-only row. That lets spend_per_minute be a plain Kafka or JDBC sink rather than an upsert sink. State is small: one partial sum per customer per open window, cleared when the window fires. Flink windowing, in depth covers triggers, lateness and session windows.

Worked example, part 2: the velocity rule in Java

The alert is a per-customer state machine with a timeout, which is what a KeyedProcessFunction is for. It counts orders in a 60-second window that opens on the customer's first order and uses an event-time timer to reset.

public class VelocityAlert extends KeyedProcessFunction<String, Order, Alert> {
    private transient ValueState<Long> count;
    private transient ValueState<Long> windowEnd;

    @Override
    public void open(OpenContext ctx) {
        count = getRuntimeContext().getState(new ValueStateDescriptor<>("count", Types.LONG));
        windowEnd = getRuntimeContext().getState(new ValueStateDescriptor<>("window-end", Types.LONG));
    }

    @Override
    public void processElement(Order o, Context ctx, Collector<Alert> out) throws Exception {
        Long end = windowEnd.value();
        if (end == null) {                       // first order for this key: open a window
            end = o.ts + 60_000;
            windowEnd.update(end);
            ctx.timerService().registerEventTimeTimer(end);
        }
        long n = (count.value() == null ? 0 : count.value()) + 1;
        count.update(n);
        if (n == 6) out.collect(new Alert(o.customerId, n, end));   // alert once per window
    }

    @Override
    public void onTimer(long t, OnTimerContext ctx, Collector<Alert> out) {
        count.clear();                           // clear state, or it lives forever
        windowEnd.clear();
    }
}

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
KafkaSource<Order> source = KafkaSource.<Order>builder()
    .setBootstrapServers("kafka:9092")
    .setTopics("orders")
    .setGroupId("velocity-alerts")
    .setStartingOffsets(OffsetsInitializer.earliest())
    .setValueOnlyDeserializer(new OrderDeserializer())
    .build();
WatermarkStrategy<Order> wm = WatermarkStrategy
    .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
    .withTimestampAssigner((o, recordTs) -> o.ts)
    .withIdleness(Duration.ofMinutes(1));

env.fromSource(source, wm, "orders").uid("orders-source")
   .keyBy(o -> o.customerId)
   .process(new VelocityAlert()).uid("velocity-rule").name("velocity-rule")
   .sinkTo(alertSink).uid("alerts-sink");
env.execute("velocity-alerts");

Trace one customer. Orders arrive at 12:00:03, :09, :15, :22, :31 and :40 event time. The first opens a window ending at 12:01:03 and registers a timer. The sixth makes the count 6 and emits one alert. When the watermark passes 12:01:03, which happens once the source has seen an event stamped later than 12:01:08, the timer fires and clears both values. If the job crashes at :35 and the last completed checkpoint was taken after the :15 order, Flink restores a count of 3 and the pending timer, rewinds Kafka to the offsets in that checkpoint, and replays the :22 and :31 orders to rebuild the count of 5, so the alert still fires once in state terms.

Three details in this code are deliberate. OffsetsInitializer.earliest() applies only on a fresh start; on restore, offsets come from the checkpoint. withIdleness stops an empty partition from holding the watermark back forever. And every stateful operator has a uid, which matters in the next section.

Running it: deployment modes and configuration

Flink runs a job in one of two ways in current releases. In application mode the cluster exists for one application: the JobManager runs your main(), and the cluster goes away with the job. Prefer it in production: one bad job cannot take down another. In session mode a long-running cluster accepts many jobs, which saves start-up time for short or interactive work such as a SQL gateway, at the cost of shared failure domains. On Kubernetes, the Flink Kubernetes Operator manages either mode through custom resources and handles upgrades with savepoints.

Flink 2.x reads only conf/config.yaml (the legacy flink-conf.yaml format is gone). A sensible starting point for the job above:

execution.checkpointing.interval: 60s
execution.checkpointing.dir: s3://payments-flink/checkpoints/velocity
state.backend.type: rocksdb
pipeline.max-parallelism: 720
pipeline.generic-types: false
restart-strategy.type: exponential-delay

RocksDB keeps state on local disk, so state size is bounded by disk rather than heap; Flink state backends compared explains when the heap backend or the disaggregated ForSt backend in 2.x is the better fit. Maximum parallelism fixes the number of key groups, the unit in which keyed state is redistributed when you rescale, and cannot be changed later without losing state. Pick a value with many divisors, well above the parallelism you expect.

Keeping a job upgradable

A streaming job runs for months, and you will change its code. The upgrade path is: stop the job with a savepoint, deploy the new jar, start it from the savepoint. Whether that works is decided by choices made on day one.

  • Set a uid on every stateful operator and on sources and sinks. Without one, Flink derives an id from the graph's shape, so inserting a single map operator can change ids and the savepoint's state no longer matches. Restore then fails, or, if you allow non-restored state, silently drops it.
  • Keep state types evolvable. POJOs and Avro types support adding and removing fields. Types Flink cannot analyse fall back to Kryo, which is slow and breaks on many class changes. pipeline.generic-types: false turns the fallback into a start-up error, so you find out in development.
  • Name state descriptors once and never reuse a name for a different type.
  • Test the restore. Take a savepoint of the old version in staging and start the new version from it before every release. Flink savepoints architecture covers formats, ownership and rescaling.

Moving to Flink 2.x

Flink 2.0 was a breaking release. The DataSet API, the legacy SourceFunction and SinkFunction interfaces and StreamExecutionEnvironment.addSource() are gone; sources use the unified Source API through env.fromSource(source, watermarkStrategy, name) and sinks use sinkTo. Rich functions use open(OpenContext) instead of open(Configuration), and the configuration file is config.yaml. The 1.20 line, the last 1.x minor release, still receives patch releases, which gives teams time to migrate.

Migrate in this order: move the job to the unified Source and Sink connectors while still on 1.20, remove any DataSet code, convert the configuration file, then switch the runtime. Check that each connector you use has a release built for 2.x before you plan the date, because connectors are released separately from Flink itself.

Failure modes you will meet early

  • Output stops while input flows. An idle or empty Kafka partition holds the watermark back, so no window or timer fires. Configure idleness, and alert on watermark lag per source.
  • State grows without bound. A process function that never clears state, or a SQL regular join without a state TTL, keeps every key forever. Clear state in timers and set table.exec.state.ttl for SQL joins after deciding how stale a join partner may be.
  • Checkpoints time out under backpressure. A slow sink or a hot key fills buffers, markers queue behind data and checkpoints expire. Find the bottleneck operator in the web UI's backpressure view before raising the timeout.
  • Hot keys. keyBy sends one key to one instance; a key with a large share of the traffic saturates one slot while the rest sit idle. Pre-aggregate with a salted key, then combine.
  • Late events disappear. Events behind the watermark are dropped by windows by default. Measure how late they arrive, then choose allowed lateness or a side output.

Trade-offs and when not to use Flink

Flink is the right tool when correctness depends on event time and large keyed state: fraud rules, sessionisation, real-time feature computation, stream joins and CDC pipelines. It costs real operations work: a cluster or operator to run, checkpoint storage, state that must survive upgrades, and on-call people who understand watermarks.

If you only need to transform or route messages between Kafka topics with small state, Kafka Streams runs inside your service with no cluster. If minute-level latency is fine, a micro-batch job or an incremental warehouse model is simpler to debug.

What to do next

  1. Run the SQL example against a local Kafka with the Flink SQL client and watch a window fire as the watermark advances.
  2. Write the velocity rule as a KeyedProcessFunction and test it with Flink's operator test harness, advancing the watermark by hand.
  3. Put a uid on every stateful operator in your existing jobs and set pipeline.generic-types to false in a test run.
  4. Choose maximum parallelism and checkpoint storage before the first production deploy, and write both down.
  5. Rehearse an upgrade: savepoint, change the code, restore, and compare output against the old version.
  6. If you are on 1.x, list every legacy SourceFunction, SinkFunction and DataSet use as the 2.x migration backlog.
Key takeaway: Flink keeps per-key state next to the computation, orders work by event time with watermarks, and snapshots state and source offsets together so a crash costs a replay rather than a wrong answer. Start with SQL for aggregations and joins, and move to DataStream and process functions for per-key state machines. Most long-term pain comes from day-one choices: stable uids, evolvable state types, a sensible maximum parallelism and rehearsed savepoint restores.