In Flink SQL, how do you write a CREATE TABLE that reads JSON page views from Kafka and declares event time with a watermark?
answer
- options live in the WITH clause
- connector, topic, servers, format
- rowtime must be a millisecond timestamp
- computed column converts epoch millis
- WATERMARK FOR col AS col minus delay
basics
~20 sDeclare the payload columns, an event-time column of type TIMESTAMP(3) or TIMESTAMP_LTZ(3), computed if needed, and WATERMARK FOR it AS that column minus an allowed delay; the WITH clause sets 'connector' = 'kafka', topic, bootstrap servers and 'format' = 'json'.
solid answer
~40 sA Flink `CREATE TABLE` stores no data; it registers a schema plus connector options in the current catalog. Physical columns mirror the JSON fields. If the payload carries epoch milliseconds, a computed column such as `event_time AS TO_TIMESTAMP_LTZ(view_ts_ms, 3)` yields a valid time attribute, and `WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND` makes it the table's event-time attribute with a five-second out-of-orderness bound. Kafka's own record fields arrive as metadata columns, for example `kafka_ts TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' VIRTUAL`, where `VIRTUAL` keeps a column out of the schema used when the table is written to. The `WITH` clause carries `'connector' = 'kafka'`, `'topic'`, `'properties.bootstrap.servers'`, `'format' = 'json'`, and usually `'properties.group.id'` and `'scan.startup.mode'`, whose default is `group-offsets`.
code
sql · 15 linesCREATE TABLE page_views (
user_id STRING,
page_url STRING,
view_ts_ms BIGINT,
event_time AS TO_TIMESTAMP_LTZ(view_ts_ms, 3),
kafka_ts TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' VIRTUAL,
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'page-views',
'properties.bootstrap.servers' = 'broker:9092',
'properties.group.id' = 'page-view-counter',
'scan.startup.mode' = 'earliest-offset',
'format' = 'json'
);go deeper
Be able to write this DDL from memory: physical columns, a computed timestamp, the WATERMARK FOR clause, and the four or five WITH options a Kafka source needs.
Explain the difference between physical, computed and metadata columns, why VIRTUAL exists, and how the watermark expression is evaluated per row but emitted periodically.
Show you choose the watermark bound from measured disorder, pin scan.startup.mode deliberately for replays, and know which Kafka client properties Flink overrides.
Discuss who owns shared table definitions across teams, how option defaults become hidden contracts, and why DDL belongs in reviewed, versioned catalogs rather than ad hoc sessions.
## What a CREATE TABLE does in Flink In Flink SQL, `CREATE TABLE` does not create storage. It registers a **table definition** in the current catalog: a schema, optional time attributes and constraints, and a `WITH` clause of **connector options** that tell Flink how to reach the data. The rows stay in the Kafka topic; Flink reads them only when a query that uses the table is submitted. The same definition can serve as a source (`SELECT ... FROM page_views`) and as a sink (`INSERT INTO page_views ...`), which is why some columns are marked as read-only. ## The three kinds of column | Kind | Example | Read by SELECT | Written by INSERT INTO | |---|---|---|---| | Physical | `user_id STRING` | yes | yes | | Computed | `event_time AS TO_TIMESTAMP_LTZ(view_ts_ms, 3)` | yes | no | | Metadata | `kafka_ts TIMESTAMP_LTZ(3) METADATA FROM 'timestamp'` | yes | yes, unless `VIRTUAL` | - **Physical columns** map to fields of the record value that the `format` decodes; with `'format' = 'json'` they are matched by field name. - **Computed columns** are expressions over other columns of the same table. The planner turns them into a projection after the source, and they cannot be the target of an `INSERT INTO`. They are the usual way to turn an epoch-millisecond `BIGINT` into a timestamp. - **Metadata columns** expose what the connector knows about each record rather than the payload. The Kafka connector offers keys such as `timestamp` (`TIMESTAMP_LTZ(3)`, readable and writable), `partition` and `offset` (read-only). Adding `VIRTUAL` excludes the column from the schema used when the table is a sink, which read-only metadata needs. ## The WATERMARK clause `WATERMARK FOR rowtime_column AS watermark_strategy_expression` does two things at once: it marks the column as the table's **event-time attribute**, and it defines how the source generates watermarks. The column must be a top-level `TIMESTAMP(3)` or `TIMESTAMP_LTZ(3)` column, and it may be a computed column. The expression is evaluated for every row, and the largest value generated so far is emitted periodically, every `pipeline.auto-watermark-interval` (200 ms by default); a watermark that is null or not larger than the previous one is not emitted. The Flink documentation names three common strategies: 1. **Strictly ascending**: `WATERMARK FOR ts AS ts` - the watermark equals the largest timestamp seen so far. 2. **Ascending**: `WATERMARK FOR ts AS ts - INTERVAL '0.001' SECOND` - rows equal to the largest timestamp are not late. 3. **Bounded out-of-orderness**: `WATERMARK FOR ts AS ts - INTERVAL '5' SECOND` - tolerates rows up to five seconds behind the largest one seen, the usual choice for a Kafka topic. A table without a watermark can still be queried, but it has no event-time attribute for time-based operators to use. A table can instead declare a **processing-time attribute** with a computed column `proc_time AS PROCTIME()`, which needs no watermark. ## The WITH clause for Kafka For `flink-connector-kafka` v5.0.0 the options a page-view source typically sets are: - `'connector' = 'kafka'` - the factory identifier. - `'topic'` - one topic, or a semicolon-separated list for a source. - `'properties.bootstrap.servers'` - Kafka client settings are passed through with the `properties.` prefix removed; a few, such as `auto.offset.reset`, are overridden by Flink and cannot be set this way. - `'properties.group.id'` - the consumer group; optional for a source. - `'scan.startup.mode'` - `earliest-offset`, `latest-offset`, `group-offsets`, `timestamp` or `specific-offsets`; the default is `group-offsets`, which resumes from the group's committed offsets. - `'format' = 'json'` - how the record value is decoded; `'key.format'` and `'value.format'` are used instead when the key must be decoded too. ## A complete example ```sql CREATE TABLE page_views ( user_id STRING, page_url STRING, view_ts_ms BIGINT, event_time AS TO_TIMESTAMP_LTZ(view_ts_ms, 3), kafka_ts TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' VIRTUAL, `partition` INT METADATA VIRTUAL, WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'page-views', 'properties.bootstrap.servers' = 'broker:9092', 'properties.group.id' = 'page-view-counter', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' ); ``` `partition` is back-quoted because it is a reserved word. `event_time` is computed, so it is readable but never written; both metadata columns are `VIRTUAL`, so an `INSERT INTO page_views` would supply only the three physical columns. ## Mistakes that show up in review - Pointing the watermark at a `BIGINT` epoch column: the rowtime column must be a timestamp, so a computed column converts it first. - Copying the pre-1.11 property style (`'connector.type' = 'kafka'`) from old examples; Flink 2.3 documents only the `'connector' = 'kafka'` style. - Forgetting `VIRTUAL` on read-only metadata such as `offset`, which makes writing to the same table try to set a key the connector can only read. - Setting the watermark delay far larger than the real disorder: every result that waits for event time is held back by the full bound. - Assuming `scan.startup.mode` defaults to the earliest offset; the default is `group-offsets`. - Expecting the DDL to validate the topic: options are checked when a query using the table is planned, not when the definition is registered.
- How would you give the same Flink table a processing-time attribute instead of event time?Add a computed column `proc_time AS PROCTIME()` and drop the `WATERMARK` clause. Processing time follows the clock of the machine evaluating the operator, so it needs no watermark, but results then depend on when rows are processed: replaying the same topic can group rows differently from the first run.
- What does Flink's json format do with a record that has a missing field or cannot be parsed?By default `json.fail-on-missing-field` is `false`, so a missing field is read as null. A record that cannot be parsed fails the job unless `'json.ignore-parse-errors' = 'true'`, which skips bad rows and sets fields with parse errors to null. Skipping hides data loss, so pair it with a check on row counts.
- Why mark the Kafka offset metadata column as VIRTUAL?Without `VIRTUAL` the planner treats a metadata column as both readable and writable. `offset` can only be read, so the column must be excluded from the schema used when the table is an `INSERT INTO` target; `VIRTUAL` does exactly that while keeping it available to `SELECT`.
saying these in an interview costs you the question
- The WATERMARK clause drops rows older than the delay
- Any BIGINT column can be the rowtime column directly
- CREATE TABLE copies the topic's data into Flink
- scan.startup.mode defaults to reading from the earliest offset
- Computed columns can be written by INSERT INTO like physical ones