skip to content

In Flink SQL, how do you write a CREATE TABLE that reads JSON page views from Kafka and declares event time with a watermark?

level: juniorimportance: must knowfreq 58%

answer

  1. options live in the WITH clause
  2. connector, topic, servers, format
  3. rowtime must be a millisecond timestamp
  4. computed column converts epoch millis
  5. WATERMARK FOR col AS col minus delay

basics

~20 s

Declare 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 s

A 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 lines
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,
  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

for a junior

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.

for a middle

Explain the difference between physical, computed and metadata columns, why VIRTUAL exists, and how the watermark expression is evaluated per row but emitted periodically.

for a senior

Show you choose the watermark bound from measured disorder, pin scan.startup.mode deliberately for replays, and know which Kafka client properties Flink overrides.

for a principal

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