In Flink, when do you use StreamTableEnvironment.toDataStream versus toChangelogStream, and how do event time and watermarks cross between the Table and DataStream APIs?
answer
- one environment, two APIs
- insert-only versus updating tables
- RowKind on every Row
- SOURCE_WATERMARK() reuses upstream watermarks
- execute() still runs the job
basics
~20 stoDataStream 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 linesStreamExecutionEnvironment 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
Know that toDataStream is for insert-only tables and toChangelogStream is for updating ones, and that env.execute() still runs the job.
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.
Decide where the SQL/DataStream boundary sits in a real pipeline, keep it narrow, and replace deprecated toRetractStream code during upgrades.
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