In the Flink SQL client, why does a table you defined with CREATE TABLE disappear after the client restarts, and how do you keep it?
answer
- where does the definition live
- default_catalog.default_database
- in-memory, per session
- a persistent catalog stores tables
- catalog store keeps catalogs, not tables
basics
~20 sBy default the table lands in default_catalog, a GenericInMemoryCatalog whose objects live only for the session. To keep it, create the table in a persistent catalog such as a HiveCatalog, or re-run the DDL from an init script at startup.
solid answer
~40 sEvery Flink table lives at `catalog.database.table`. Without a `USE CATALOG`, the SQL client uses `default_catalog.default_database`, backed by a `GenericInMemoryCatalog`, and all its objects exist only for the lifetime of the session, so a restart forgets the definition. The data is untouched - it was always in Kafka - and a job already submitted keeps running because it carries its compiled plan. To keep definitions, register a persistent catalog, for example `CREATE CATALOG hive_catalog WITH ('type' = 'hive', ...)`, `USE CATALOG hive_catalog`, then create the table there. A **catalog store** (`table.catalog-store.kind` = `file`) persists only catalog configurations, not tables. A `JdbcCatalog` exposes existing database tables read-only and cannot store Flink's DDL. `CREATE TEMPORARY TABLE` is session-only in any catalog.
code
sql · 18 linesSHOW CURRENT CATALOG; -- default_catalog: in-memory, lost on restart
CREATE CATALOG hive_catalog WITH (
'type' = 'hive',
'hive-conf-dir' = '/opt/hive-conf'
);
USE CATALOG hive_catalog;
CREATE TABLE page_views (
user_id STRING,
page_url STRING,
view_ts_ms BIGINT
) WITH (
'connector' = 'kafka',
'topic' = 'page-views',
'properties.bootstrap.servers' = 'broker:9092',
'format' = 'json'
);go deeper
Know that default_catalog is in-memory and per session, and that tables meant to last belong in a persistent catalog.
Distinguish a persistent catalog from a catalog store and an init script, and explain temporary-table shadowing.
Set up catalogs so production definitions are shared and reviewed, and know that a lost definition does not stop jobs already running.
Decide where table definitions are the source of truth across teams - metastore, repository or both - and how changes to them are reviewed and rolled out.
## Where a table definition lives In Flink SQL, every table, view and function is addressed as `catalog.database.object`. A **catalog** holds metadata: table schemas, connector options, views, functions. It holds no rows; the rows stay in Kafka, a database or files. When a statement names only `page_views`, Flink resolves it against the current catalog and database. Out of the box a session has one catalog, `default_catalog` (the name comes from `table.builtin-catalog-name`), with one database, `default_database`. Unless you run `USE CATALOG`, that is where `CREATE TABLE` writes. ## Why it disappears `default_catalog` is a **`GenericInMemoryCatalog`**: all objects in it are available only for the lifetime of the session. Closing the SQL client ends the session and discards the definitions. Nothing is wrong with the table or its data; only the metadata is gone. Running `SHOW CURRENT CATALOG` before creating tables is a quick way to see where they will land. A job already submitted with `INSERT INTO ... SELECT ...` keeps running on the cluster after the client exits: its plan was compiled and submitted with the table definitions resolved, and it does not look them up again. What is lost is the ability to run new statements against `page_views` without re-creating it. ## Keeping definitions 1. **Use a persistent catalog.** A `HiveCatalog` stores Flink table definitions in a Hive Metastore, where every session connected to that catalog can see them until they are dropped. It needs the Hive connector dependencies on the classpath. 2. **Persist the catalog registration itself.** A **catalog store** saves the configuration of catalogs created in a session, so they are registered again when a new session starts. The default kind is `generic_in_memory`; `table.catalog-store.kind` = `file` with `table.catalog-store.file.path` selects `FileCatalogStore`, one YAML file per catalog. 3. **Replay DDL at startup.** `sql-client.sh -i init.sql` runs an initialization script, so the in-memory catalog is refilled each time. This works without extra infrastructure, but the source of truth becomes a file that every user must share. | Option | What survives a restart | What it needs | |---|---|---| | `GenericInMemoryCatalog` (default) | nothing | nothing | | `HiveCatalog` | table, view and function definitions | a Hive Metastore | | `JdbcCatalog` | nothing Flink creates; it reads existing DB tables | a relational database the JDBC connector ships a catalog for | | Catalog store (`file`) | catalog registrations, not tables | a file-system path | | `-i init.sql` | whatever the script re-creates | a shared script | The `JdbcCatalog` row surprises people: in `flink-connector-jdbc` v4.1.0 its implementations support only reading methods such as `listTables` and `getTable`. It maps existing database tables into Flink so you do not write their DDL by hand, but it cannot store a Kafka table you define. ## Temporary tables and shadowing `CREATE TEMPORARY TABLE` always creates a session-scoped table, even when the current catalog is persistent; it is never written to the catalog. If a temporary table has the same identifier as a permanent one, it **shadows** it: queries in that session use the temporary table until it is dropped. This is handy for testing a query against a sample, and confusing when forgotten. ## Doing the same from code In a Table API program the same choices apply. `tableEnv.createCatalog(name, CatalogDescriptor.of(name, configuration))` creates a catalog from its configuration and records the descriptor in the catalog store, so a later session with the same store gets the catalog back. The older `registerCatalog(name, catalogInstance)` still exists in 2.3 but is deprecated in favour of `createCatalog`, which is the one that records the descriptor in the store. `tableEnv.useCatalog(name)` and `useDatabase(name)` then play the role of `USE CATALOG` and `USE`. ## Common confusions - Believing `CREATE TABLE` copies data into Flink; it registers metadata only. - Expecting the catalog store to keep tables; it keeps catalog configurations. - Assuming a running job stops when its catalog entry disappears; it does not. - Registering catalogs in code with the deprecated `registerCatalog`; in 2.3 `createCatalog(name, CatalogDescriptor)` is preferred because it also records the descriptor in the catalog store.
- Can a Flink JdbcCatalog hold the Kafka page_views table you create?No. In flink-connector-jdbc v4.1.0 the JDBC catalogs support only read methods such as listing databases and tables and getting a table. They expose existing database tables to Flink; creating a Kafka table in them is not supported.
- What happens to an INSERT INTO job already submitted when the Flink SQL client restarts?It keeps running. The statement was planned with its table definitions resolved and submitted as a job; the job never consults the catalog again. Only new statements that reference the forgotten table fail until it is defined again.
saying these in an interview costs you the question
- CREATE TABLE copies the topic's data into Flink storage
- The file catalog store keeps every table definition
- Restarting the SQL client stops the jobs it submitted
- A JdbcCatalog can store any table you create in Flink
- CREATE TEMPORARY TABLE in a Hive catalog is persisted