skip to content

In Flink, when do you use StreamTableEnvironment.toDataStream versus toChangelogStream, and how do event time and watermarks cross between the Table and DataStream APIs?

level: middleimportance: should knowfreq 40%

answer

  1. one environment, two APIs
  2. insert-only versus updating tables
  3. RowKind on every Row
  4. SOURCE_WATERMARK() reuses upstream watermarks
  5. execute() still runs the job

basics

~20 s

toDataStream accepts only insert-only tables; an updating table such as a GROUP BY result needs toChangelogStream, whose Rows carry a RowKind. Inbound, a Schema with SOURCE_WATERMARK() reuses the stream's watermarks; outbound, a single rowtime column becomes the record timestamp.

solid answer

~40 s

`StreamTableEnvironment` bridges the two APIs over one `StreamExecutionEnvironment`. `fromDataStream(stream)` or `createTemporaryView(name, stream)` exposes an insert-only stream as a table. To keep event time, pass a `Schema` with `columnByMetadata("rowtime", "TIMESTAMP_LTZ(3)")` and `watermark("rowtime", "SOURCE_WATERMARK()")`, which reads each record's timestamp and reuses the watermarks already generated upstream. Going back, `toDataStream(table)` works only for insert-only tables; on an updating table it fails with `doesn't support consuming update changes`. `toChangelogStream(table)` returns `DataStream<Row>` where `row.getKind()` gives `INSERT`, `UPDATE_BEFORE`, `UPDATE_AFTER` or `DELETE`, and an overload takes `ChangelogMode.upsert()` to drop `UPDATE_BEFORE`. A single rowtime column becomes the stream record's timestamp and watermarks propagate. The older `toAppendStream` and `toRetractStream` still exist in 2.3 but are deprecated.

code

java · 25 lines
java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

// PageView is a POJO with public fields userId, pageUrl, viewTsMs;
// timestamps and watermarks were assigned upstream with a WatermarkStrategy.
DataStream<PageView> viewStream = buildPageViewStream(env);

Table views = tableEnv.fromDataStream(
    viewStream,
    Schema.newBuilder()
        .columnByMetadata("rowtime", "TIMESTAMP_LTZ(3)")
        .watermark("rowtime", "SOURCE_WATERMARK()")
        .build());
tableEnv.createTemporaryView("page_views", views);

Table counts = tableEnv.sqlQuery(
    "SELECT userId, COUNT(*) AS views FROM page_views GROUP BY userId");

// toDataStream(counts) would fail: the table is updating.
DataStream<Row> changes = tableEnv.toChangelogStream(counts);
changes
    .filter(row -> row.getKind() != RowKind.UPDATE_BEFORE)
    .print();

env.execute("page-view-counts");

go deeper

for a junior

Know that toDataStream is for insert-only tables and toChangelogStream is for updating ones, and that env.execute() still runs the job.

for a middle

Explain RowKind on each Row, how SOURCE_WATERMARK() and a metadata rowtime column carry event time in, and how a single rowtime column carries it out.

for a senior

Decide where the SQL/DataStream boundary sits in a real pipeline, keep it narrow, and replace deprecated toRetractStream code during upgrades.

for a principal

Set team guidance on when logic belongs in SQL, in Process Table Functions or in DataStream code, trading optimiser benefits against flexibility and maintenance.

## Why cross between the two APIs Flink SQL and the Table API express relational logic that the planner optimises; the DataStream API expresses record-at-a-time logic with explicit state and timers. Real pipelines often need both: a DataStream source with no table connector feeding a SQL aggregation, or a SQL result feeding custom processing. `StreamTableEnvironment.create(env)` builds a table environment on top of a `StreamExecutionEnvironment`, so both halves compile into **one job graph**. ## From DataStream to Table - `fromDataStream(DataStream<T>)` interprets an **insert-only** stream as a table and derives the columns from the type, for example the public fields of a POJO. - `fromDataStream(DataStream<T>, Schema)` does the same with an explicit schema: extra computed columns, metadata columns, a watermark. - `createTemporaryView(String, DataStream)` registers the stream under a name for SQL; it is a shortcut for `createTemporaryView(name, fromDataStream(stream))`. - `fromChangelogStream(DataStream<Row>)` interprets rows that already carry a `RowKind` as an updating table; an overload takes a `ChangelogMode`, such as `ChangelogMode.upsert()`. ## Event time across the boundary Into a table, the DataStream records already have timestamps and the stream already has watermarks. A schema can adopt both: ```java Table views = tableEnv.fromDataStream( viewStream, Schema.newBuilder() .columnByMetadata("rowtime", "TIMESTAMP_LTZ(3)") .watermark("rowtime", "SOURCE_WATERMARK()") .build()); ``` `columnByMetadata("rowtime", ...)` exposes the stream record's timestamp as a column, and `SOURCE_WATERMARK()` tells the table to reuse the upstream watermarks instead of generating new ones. A schema can instead compute a rowtime column from a field with `columnByExpression` and declare its own strategy, such as `rowtime - INTERVAL '10' SECOND`. Out of a table, if the result has a **single rowtime column**, `toDataStream` and `toChangelogStream` write it into each stream record's timestamp and propagate the watermarks, so downstream DataStream event-time operators keep working. ## From Table to DataStream | Method | Accepts | Emits | Status in 2.3 | |---|---|---|---| | `toDataStream(Table)` | insert-only tables | `Row`, or a requested class via an overload | current | | `toChangelogStream(Table)` | insert-only or updating | `Row` with `RowKind` set | current | | `toChangelogStream(Table, Schema, ChangelogMode)` | updating, limited to a mode | `Row` with `RowKind` set | current | | `toAppendStream` | insert-only | typed records | deprecated | | `toRetractStream` | updating | `Tuple2<Boolean, T>` | deprecated | Calling `toDataStream` on a `GROUP BY` result fails at planning with `Table sink 'Unregistered_DataStream_Sink_1' doesn't support consuming update changes`. The fix is either `toChangelogStream`, whose rows are inspected with `row.getKind()`, or a query that produces an insert-only result. `registerDataStream`, an older way to name a stream, was removed in Flink 2.0. ## Execution: who runs the job Each `toDataStream` or `toChangelogStream` call compiles the table sub-pipeline and adds it to the DataStream pipeline builder. The job then runs on `env.execute()` (or `executeAndCollect()`); a Table API execution such as `executeSql(...)` does not trigger these DataStream parts. Forgetting `execute()` after a conversion is a common reason for a program that builds a plan and produces nothing. ## Mistakes that show up in review - Calling `toDataStream` on an aggregation and then working around the planning error by rewriting the query, when `toChangelogStream` was the right call. - Filtering out `UPDATE_BEFORE` rows downstream and then feeding the remaining rows into a DataStream aggregation that sums them: every update is then counted as a new value unless the consumer overwrites by key. - Defining a fresh watermark strategy in the `Schema` when the DataStream already has one, so the table and the stream disagree about event-time progress; `SOURCE_WATERMARK()` avoids that. - Keeping old `toRetractStream` code through an upgrade because it still compiles; it is deprecated, and its `Tuple2<Boolean, T>` output does not distinguish an update from a delete followed by an insert. ## When to stay in SQL and when to drop down - **Stay in SQL** for projections, filters, aggregations and anything the planner can reorder, push down or reuse; it also keeps the logic readable to more of the team. - **Drop to DataStream** for logic that needs custom keyed state with its own timers, side outputs, or a source or sink that has only a DataStream connector. - **Consider Process Table Functions**, available since Flink 2.1, before dropping down: they bring state and timers into SQL for some of these cases. - **Keep the boundary narrow**: every crossing converts types, and schema changes on either side must be mirrored on the other.

  • When should a Flink job drop from SQL to the DataStream API rather than stay in SQL?
    When the logic needs custom keyed state with its own timers, side outputs, or a connector that exists only for DataStream. Relational steps stay in SQL, where the planner optimises them. Since Flink 2.1, Process Table Functions cover some state-and-timer cases without leaving SQL.
  • What does toChangelogStream(table, schema, ChangelogMode.upsert()) change about the output?
    It asks for an upsert changelog, so the stream carries no `UPDATE_BEFORE` rows: consumers overwrite by key on `UPDATE_AFTER` and remove on `DELETE`. The method throws if the updating table cannot be represented in the requested mode, so the result must have a key consumers can upsert on.
  • Why does a Flink program that calls toDataStream(...).print() but never env.execute() output nothing?
    The conversion only adds the compiled table pipeline to the DataStream graph. Nothing runs until `env.execute()` or `executeAndCollect()` submits that graph; a Table API execution does not trigger the DataStream branches.

saying these in an interview costs you the question

  • toDataStream works for any table, updates included
  • toRetractStream was removed in Flink 2.0
  • Converting a table to a stream loses its event time
  • Every conversion starts a separate Flink job
  • SOURCE_WATERMARK() generates fresh watermarks from rowtime