The Kafka Connect Iceberg sink reports every task RUNNING and the Flink job writing the table completes its checkpoints. Then an analyst queries the table in Spark or Trino and gets no new rows, or a downstream job fails on startup with NoSuchTableException or NoSuchNamespaceException. The writer is committing, but to a table the reader is not looking at.
A catalog mismatch is a writer and a reader resolving the same table name through different Iceberg catalogs, warehouse paths, catalog names or namespaces, so each side loads a different table or none at all. The Iceberg table spec says a commit swaps the table’s metadata pointer, and that the pointer can be stored in a metastore or database. The Spark register_table docs warn that having the same metadata.json registered in more than one catalog can lead to missing updates, loss of data and table corruption, which describes what can happen when two pointers to one table move apart.
This page covers the case where the writer succeeds and the reader cannot see its work. For a near miss, go to the page that fits.
- Writer commits fail with
CommitFailedException: Iceberg commits failing or conflicting. - The table is found but a new field never arrives: Kafka to Iceberg pipeline breaks after a schema change.
- Data lands but queries crawl over thousands of tiny files: too many small files and how to compact them.
For how catalogs, snapshots and metadata files fit together, see the complete Iceberg guide.
Rule out a writer that is not writing
Confirm the writer is really committing before comparing catalogs. Two signals show whether a sink is committing: a task error and a consumer that has stopped keeping up. The Kafka Connect REST API has a status endpoint that reports whether the connector is running, failed or paused, the error information if it has failed, and the state of every task.
GET /connectors/{name}/status
The Kafka consumer fetch metrics include the maximum lag in records for any partition, which shows whether the sink’s consumer is keeping up with the topic.
records-lag-max
A failed task or a growing lag is a writer problem, and the Kafka Connect troubleshooting guide is the place to start. If every task is RUNNING and lag stays flat, the writer is committing somewhere and the rest of this page applies.
Rule out a reader that has not refreshed
A reader can be pointed at the right table and still miss recent commits. The Iceberg spec says readers use the snapshot that was current when they loaded the table metadata, and are not affected by changes until they refresh and pick up a new metadata location.
Engines also cache catalog entries. The Iceberg catalog properties list caching as on by default with a 30 second expiry, and the Spark configuration gives the same 30 second default. These are the catalog and Spark defaults.
cache-enabled=true
cache.expiration-interval-ms=30000
The Flink configuration gives a different default for Flink catalogs, where a negative value disables expiration.
cache.expiration-interval-ms=-1
If a new session or a restarted job sees the rows, the reader was holding an older snapshot, not looking at a different table. If a fresh reader still sees nothing, or fails with table not found, compare the catalogs.
Compare the catalog settings on both sides
Every Iceberg engine is told which catalog to use through its own settings, and nothing forces two engines to agree. Compare each row below on the writer and on the reader. The first row that differs is the mismatch.
The catalog name is only a label. Two engines can use the same name for different catalogs, or different names for the same one, so compare the properties under the name, not the name. The sink docs give the name setting a default of iceberg, and a Spark catalog can be called hive_prod or anything else.
iceberg.catalog=iceberg
Check each setting in turn
The Apache Iceberg Kafka Connect sink is the self-managed route from a Kafka topic into an Iceberg table, so its catalog block is the first thing to compare. The sink docs say the properties that start with iceberg.catalog. are required for connecting to the Iceberg catalog, and that the core catalog types in the default distribution include REST, Glue, DynamoDB, Hadoop, Nessie, JDBC, Hive and BigQuery Metastore. They add that the Hive catalog needs the distribution that includes the Hive metastore client.
Catalog type and endpoint
The sink picks its catalog type with one property, or names a catalog class for the types that have no short value. The sink docs accept rest, hive or hadoop for the type.
iceberg.catalog.type=rest
iceberg.catalog.catalog-impl=<catalog class>
A Spark catalog sets the type under its own name, and the Spark docs list hive, hadoop, rest, glue, jdbc and nessie as values.
spark.sql.catalog.<name>.type=rest
A Flink catalog uses a different key in its CREATE CATALOG statement. The Flink docs list hive, hadoop and rest as values and say the Glue, JDBC and Nessie catalogs use the class name instead.
catalog-type=rest
catalog-impl=<catalog class>
A sink on a Hive metastore and a reader on a REST catalog are two different catalogs, even when the REST catalog serves the same storage. The endpoint also has to match, because two REST catalogs, or a staging and a production metastore, share a type and differ only in the address.
The REST catalog protocol is a common API for interacting with any Iceberg catalog, so a writer and a reader that both use rest talk to one service whatever store sits behind it. The same page says that root metadata is written by the catalog service and that the protocol supports secure table sharing through credential vending or remote signing. For a REST catalog the address is therefore the identity to compare, and a reader that reaches a different service reaches a different set of tables.
iceberg.catalog.uri=https://catalog-service
Warehouse path
The catalog properties define the warehouse as the root path of the data warehouse.
iceberg.catalog.warehouse=s3a://bucket/warehouse
For a Hive catalog in Flink, the Flink configuration says the warehouse value overwrites the Hive setting below when both hive-conf-dir and warehouse are set. A job that sets one and a job that relies on the other can disagree without either showing an error.
hive.metastore.warehouse.dir
For a REST catalog the warehouse can be a server-side setting. The REST catalog spec says clients first call the config route, which can take a warehouse query parameter, and that the server’s overrides can set the warehouse location. It returns Not Found when the requested warehouse does not exist.
GET /v1/config?warehouse=<warehouse>
Storage settings can belong to the same comparison. On S3 compatible storage, the sink docs say the catalog block may also need these properties, depending on the setup, so check them on the reader’s catalog as well.
iceberg.catalog.s3.endpoint
iceberg.catalog.s3.path-style-access
Namespace and current catalog
The sink names its tables with a namespace and a table name in a comma separated list.
iceberg.tables=default.events
The REST catalog spec lists NoSuchNamespaceException and NoSuchTableException as the not found errors for table routes, which is what a reader sees when the namespace or table does not exist in the catalog it asked.
A query that leaves out the catalog and namespace uses the engine’s current ones. The Spark configuration says Spark keeps track of the current catalog and namespace, and a catalog can set default-namespace. Run this on the reader to see which ones an unqualified query uses.
SHOW CURRENT NAMESPACE;
The same Spark docs say spark_catalog is Spark’s built-in session catalog, and that Iceberg’s SparkSessionCatalog wraps it so one catalog can load both Iceberg and non Iceberg tables from the same metastore. A separate catalog defined with SparkCatalog is a different catalog even when it points at that same metastore, so a name qualified through one can fail to resolve through the other.
Glue account and region
The AWS integration docs say there is a unique Glue metastore in each AWS account and each region, and the Glue catalog picks one from the default credentials and region. A sink running with one account’s credentials and a reader running with another’s see two catalogs. The catalog property below points a catalog at a Glue metastore in a different account.
glue.id=<aws-account-id>
Nessie branch
A Nessie catalog works on a branch or tag, set with the optional ref property. The Nessie docs say every Iceberg transaction becomes a Nessie commit, and that the history can be listed, merged or cherry-picked across branches. A sink that commits to a work branch and a reader on the main branch can see different table states until the branch is merged.
iceberg.catalog.ref=main
A table loaded by path
A Flink DataStream job can skip the catalog entirely. The Flink writes docs show a table loaded from a storage path with the call below, so the job reaches the table through that path and not through a catalog entry. A table that exists only at a path has no entry in another catalog, so a reader that resolves names through that catalog can report it as not found.
TableLoader.fromHadoopTable(path, hadoopConf)
A warehouse engine with its own catalog integration
Each engine that reads the table has its own catalog configuration. Snowflake’s Iceberg table docs, for example, describe a catalog integration that connects Snowflake to an external Iceberg catalog, and say one is needed when a table is managed by AWS Glue. Iceberg keeps one copy of the data that several engines can read, so every engine added to the pipeline is one more catalog setting to compare when a reader reports table not found.
Compare the metadata location each side resolves
Configuration shows what each side was told to use. The metadata location shows the table each side actually loaded. Under the spec, a commit swaps the table’s metadata pointer to a new metadata file, so the writer and the reader are on the same table only if they resolve to files in the same metadata path.
From Spark, the metadata tables list the table’s metadata files with their timestamps. Run the same query through the writer’s catalog and the reader’s catalog, and compare the newest file on each side.
SELECT timestamp, file
FROM prod.db.events.metadata_log_entries
ORDER BY timestamp DESC;
On a Hive catalog, the Hive docs show the current metadata file held in the metadata_location table property. On a REST catalog, the spec says the load table response returns the file location of the table metadata in the metadata-location field.
GET /v1/{prefix}/namespaces/{namespace}/tables/{table}
A table has a physical identity, its storage location, and a logical one, its catalog, namespace and name. The OpenLineage symlinks facet exists because a dataset can be referenced by a file path in one tool and by a table name in another, so comparing names alone can hide that two names are one table or that one name is two.
If the two sides print metadata files under different paths, there are two tables. If one side cannot load the table at all, the table exists only in the other side’s catalog. Find out which one holds the rows the pipeline has written since it started, because that is the table to keep.
Fix it, one change at a time
Change one configuration at a time and check the metadata location again after each change. Do not register the same table in a second catalog to make both sides happy.
Point the reader at the writer’s catalog
When the writer’s table holds the data, the lowest-risk fix is to change the reader. Give the Spark or Flink catalog the same type, endpoint and warehouse as the sink’s catalog block, and qualify the query with the catalog and namespace. Nothing on the table changes.
Two engines load the same tables when both are given the same catalog type and the same address. The Spark docs describe the address as the Hive metastore URL for a hive catalog and the REST URL for a rest catalog, and the Flink docs describe it as the Hive metastore’s thrift URI.
spark.sql.catalog.<name>.type=hive
spark.sql.catalog.<name>.uri=thrift://host:port
Point the writer at the reader’s catalog
When the writer has been landing data in the wrong catalog, change the sink and stop it from creating another copy. The Kafka Connect sink docs list auto-creation of destination tables as false by default. With it set to true, a sink pointed at the wrong catalog creates the destination table there, which is one way the stray copy appears.
iceberg.tables.auto-create-enabled=false
Changing the catalog does not move the rows already written to the stray table. Backfill them from the topic, or move the table as below.
Move a table between catalogs with register_table
The Spark register_table procedure creates a catalog entry for a metadata.json file that already exists. The docs warn that having the same metadata.json registered in more than one catalog can lead to missing updates, loss of data and table corruption, and say to use the procedure only when the table is no longer registered in an existing catalog, or when moving a table between catalogs.
CALL target.system.register_table(
table => 'db.events',
metadata_file => 's3://bucket/warehouse/db/events/metadata/<newest>.metadata.json'
);
Stop every writer first, take the newest metadata file from the source catalog, register it in the target, then remove the entry from the source before any writer starts again. On a Hive catalog, the Hive docs only allow changing metadata_location to a path that holds the exact same metadata JSON.
Remove the stray copy without deleting data
The Spark DDL docs say that from 0.14 DROP TABLE only removes the table from the catalog, and DROP TABLE ... PURGE also deletes the table contents. For a copy registered in a catalog that holds a pointer in a metastore or service (Hive, Glue, JDBC, REST, Nessie), use plain DROP TABLE on the copy you no longer want, and only use PURGE once the metadata location check shows its files are not shared with the table you are keeping. A Hadoop catalog table lives in its own directory, and HadoopCatalog.dropTable deletes that table directory recursively even without PURGE, so do not drop it from that catalog until the metadata location check shows the files are not shared with the table you are keeping.
DROP TABLE wrong_catalog.db.events;
Where Kpow and Flex fit
This page is published by Factor House, which makes Kpow, its Kafka management product, and Flex, its Flink job management and monitoring product, so weigh this section with that in mind.
Kafka Connect. Kpow’s Kafka Connect management documentation describes a View Config action on the Connect Explore page that opens a connector’s configuration, and says sensitive config values appear redacted. That is where a sink’s iceberg.catalog.* block can be read without opening the worker. The same page describes restarting individual tasks and viewing the stacktraces of tasks in an ERROR state. Kpow does not connect to an Iceberg catalog, so it cannot show which table a reader resolves or compare metadata locations.
Flink jobs, what Flex shows. Flex’s job inspection documentation describes a Configuration tab that shows the configuration settings used to submit the job, including the user config and the job execution config. A catalog passed to the job as configuration shows up there.
Flink jobs, what Flex does not show. A catalog defined inside the job, such as a CREATE CATALOG statement, is not part of that view. Flex connects to each Flink cluster through its REST API, as its cluster configuration documentation describes, and does not read the Iceberg catalog or the table.
For choosing a tool to watch Flink jobs, the comparison of tools for self-managed Apache Flink sets Flex beside the Flink web UI and the REST API.
FAQ
Why does an Iceberg query return no rows when the sink is committing?
The reader is either holding an older snapshot or resolving the table through a different catalog. Start a fresh session to rule out the first. If it still sees nothing, compare the catalog type, uri, warehouse and namespace on both sides, then compare the metadata location each side loads.
What causes NoSuchTableException or NoSuchNamespaceException on an Iceberg table?
The catalog the reader asked does not have that namespace or table. That usually means the reader is configured with a different catalog, account, region, branch or warehouse than the writer, or the query used the engine’s current namespace instead of the one the writer used.
Can the same Iceberg table be registered in two catalogs?
It can be, but it should not be. The Iceberg docs warn that registering the same metadata.json in more than one catalog can lead to missing updates, loss of data and table corruption. Use register_table only to move a table, and remove it from the old catalog before any writer starts.
Does the catalog name have to match between engines?
No. The name is a label each engine gives its catalog configuration. Two engines see the same table when the catalog type, endpoint, warehouse and namespace match, whatever each one calls the catalog.