CREATE EXTERNAL TABLE … ENGINE = KAFKA¶
Map a MatrixOne external table to one partition of an Apache Kafka topic using the
ENGINE = KAFKAclause. The table is read-only, each Kafka message value must parse to exactly one record, and per-message metadata is exposed through synthetic__mo_*columns.
Description¶
CREATE EXTERNAL TABLE ... ENGINE = KAFKA maps a MatrixOne external table to a Kafka topic partition. Rows are streamed from the topic when the table is queried. Each Kafka message value must parse to exactly one record: for format=csv, one CSV record with the configured separator and a field count equal to the declared column count; for format=jsonl, one JSON object whose keys include every declared column name.
The table is read-only: SELECT is supported, and writes are rejected. ALTER TABLE on a Kafka external table is rejected. Options are validated at CREATE time and require no broker; a broker is dialed only when the table is scanned.
Syntax¶
CREATE EXTERNAL TABLE [IF NOT EXISTS] [db.]table_name (
column1 type1,
column2 type2,
...
) ENGINE = KAFKA [WITH ("option" = 'value' [, "option" = 'value'] ...)];
Arguments¶
Option |
Required |
Description |
|---|---|---|
|
Yes |
Comma-separated list of |
|
Yes |
Name of the Kafka topic to read. |
|
No |
Partition number to read. Default |
|
No |
|
|
No |
Consumer group name. Defaults to |
|
No |
|
|
No |
CSV only. A single character; default |
Synthetic columns¶
Kafka external tables expose synthetic metadata columns that are hidden from SELECT * but selectable by name:
Column |
Type |
Meaning |
|---|---|---|
|
bigint |
Kafka offset of the message. |
|
timestamp(3) |
Kafka message timestamp. |
|
varchar |
Message key ( |
|
varchar |
Raw message value. |
|
bigint |
Read control: last consumed offset; reading begins at |
|
bigint |
Read control: caps the message count; |
|
bigint |
Read control: ends the read after this many seconds without a new message; |
Read controls¶
Top-level <control> = <constant> conjuncts are resolved at compile time and position the read instead of filtering rows. They are consumed from the predicate:
SELECT * FROM kt
WHERE __mo_read_start_id = 1000
AND __mo_read_size = 1000000
AND __mo_read_timeout = 10;
LAST_KAFKA_MESSAGE_ID() returns the offset of the last message a completed Kafka scan returned in this session (NULL before any scan). With autocommit=false, feeding it back as the next __mo_read_start_id gives exactly-once consumption.
Usage Notes¶
Read-only: writes into a Kafka external table are rejected.
Read isolation: the scan consumes at Kafka’s
read_committedisolation; records of aborted or still-open producer transactions never surface as rows.Validation: duplicate, unknown, or invalid options are rejected at
CREATEtime. Control values that would overflow arithmetic are rejected at compile time.Reserved names: the synthetic
__mo_*column names are reserved; a table or column cannot declare them.
Examples¶
The following example creates Kafka external tables, round-trips their DDL, and demonstrates a validation error. Everything here resolves at DDL or compile time, so no broker is needed:
DROP DATABASE IF EXISTS kafka_ext_demo;
CREATE DATABASE kafka_ext_demo;
USE kafka_ext_demo;
CREATE EXTERNAL TABLE kt (a INT, b VARCHAR(100))
ENGINE = KAFKA WITH ('brokers' = '127.0.0.1:19092', 'topic' = 'events');
SHOW CREATE TABLE kt;
CREATE EXTERNAL TABLE kt2 (a INT, b VARCHAR(100))
ENGINE = KAFKA WITH ('brokers' = 'h1:9092,h2:9092', 'topic' = 't2', 'partition' = '3', 'autocommit' = 'true', 'group' = 'g2', 'format' = 'jsonl');
SHOW CREATE TABLE kt2;
-- NULL before any completed Kafka scan
SELECT last_kafka_message_id();
-- Expected-Success: false
CREATE EXTERNAL TABLE bad (a INT) ENGINE = KAFKA WITH ('brokers' = 'nohostport', 'topic' = 't');
DROP DATABASE kafka_ext_demo;
The read controls and exactly-once chaining require a reachable broker, so they are shown as a syntax template rather than a paste-and-run script:
SELECT * FROM kt
WHERE __mo_read_start_id = 1000
AND __mo_read_size = 100000;
SELECT last_kafka_message_id(); -- e.g. 4711
SELECT * FROM kt
WHERE __mo_read_start_id = 4711; -- continues after