feat(ingestion): add DeltaLake table type and detect Delta tables in Trino and StarRocks - #34172
akashverma0786 wants to merge 6 commits into
Conversation
Delta Lake tables that Glue, Trino, StarRocks and Hive already ingest are typed External or Regular. The enum needs the value before any connector can emit it: the server rejects an unknown tableType with 400 at the sink, so detection cannot be tested end to end without it. TypeScript regenerated with json2ts-generate-all.sh, Python models with make generate (gitignored). No connector change here. Refs openmetadata-collate#5994
Trino reports a Delta catalog as connector_name delta_lake in system.metadata.catalogs -- the same query the Iceberg branch already runs, so detection costs no extra round trip and reads no _delta_log. The comparison stays exact: `delta` is PrestoDB's connector name, and matching it here would mistype tables reached through it. Known limitation, inherited from the Iceberg rule rather than introduced: the catalog's connector name types every table in the catalog, so on a metastore shared by a hive and a delta_lake catalog, a plain Hive table seen through the delta_lake catalog is typed DeltaLake with no columns. Refs openmetadata-collate#5994
An external Delta catalog reports its tables as ENGINE 'DELTALAKE', upper case -- measured on 3.2.16 as HEX 44454C54414C414B45, 9 bytes. One RELKIND_MAP entry types them; engines with no entry still fall back to Regular, so no other key is touched. ENGINE names the catalog a table is read through, not the table's own storage format: the same Hive table reports DELTALAKE through a deltalake catalog and HIVE through a hive catalog. A plain Hive table visible in a Delta catalog is therefore typed DeltaLake with no columns. That is StarRocks' own answer, and no connector-side change can tell the two apart. Refs openmetadata-collate#5994
The three mapping tests read RELKIND_MAP and restated its literal, so they could only fail if someone edited the dict in the same commit; nothing exercised StarRocksSource consuming it, leaving query_table_names_and_types with no coverage at all. One parametrized test drives the real method and covers what the three asserted plus the fallback: DELTALAKE types DeltaLake, mixed case falls through to Regular, and ICEBERG/HIVE/TABLE are undisturbed. Refs openmetadata-collate#5994
The in-app connector help is the only place a user learns which table types a connector emits. Both engines decide the type from the catalog a table is read through, not from the table, so the shared-metastore mis-typing and the partitioned-table overwrite need saying out loud. The StarRocks note also records how to reach an external catalog at all: init_command with SET CATALOG, and Database Schema left empty, since it is sent as the connection's database. Refs openmetadata-collate#5994
The unit tests mock the catalog query, and the only integration coverage ran against a hive catalog, so nothing proved a DeltaLake type survives the workflow and the API. Adds a delta_lake catalog over the metastore the harness already runs, a real Delta table written by Trino, and a MetadataWorkflow whose readback asserts tableType and the column list. With the detection branch reverted the test reports Regular, so it fails on the behaviour rather than on the fixture. Refs openmetadata-collate#5994
❌ PR checklist incompleteThis PR cannot be merged until the following are addressed on its linked issue:
The fields live on the linked issue in the Shipping project (open the issue → right sidebar → Projects). After you set them, re-run this check (or push a commit) — issue/project changes do not re-trigger it automatically. Maintainers can bypass this check by adding the |
| @pytest.fixture(scope="module") | ||
| def create_delta_table(trino_container): | ||
| engine = create_engine(make_url(trino_container.get_connection_url()).set(database="delta")) | ||
| try: | ||
| with engine.connect() as conn: | ||
| conn.execute( | ||
| text( | ||
| "CREATE SCHEMA IF NOT EXISTS delta.delta_schema WITH (location = 's3a://hive-warehouse/delta_schema')" | ||
| ) | ||
| ) | ||
| conn.execute( | ||
| text("CREATE TABLE delta.delta_schema.delta_sales (id integer, region varchar, amount double)") | ||
| ) | ||
| conn.execute( | ||
| text("INSERT INTO delta.delta_schema.delta_sales VALUES (1, 'emea', 10.5), (2, 'apac', 20.25)") |
There was a problem hiding this comment.
⚠️ Bug: Delta fixture leaks delta_schema into the shared hive catalog
The delta catalog uses the same Hive metastore as minio, the Trino container is package-scoped, and create_delta_table never removes the schema or table it creates. pytest runs test_delta_lake.py before test_metadata.py, test_profiler.py and test_profiler_sampling.py. Those modules ingest the minio catalog, and their schemaFilterPattern excludes only information_schema. So they will now find minio.delta_schema.delta_sales, which is a Delta table. Trino's hive connector refuses to read it: SHOW COLUMNS in _get_columns fails with "Cannot query Delta Lake table". That failure is counted in the workflow status, and run_workflow(..., raise_from_status=True) can then fail the later modules. Whether they pass depends on test order, and running the new test alone (as the PR did) does not catch it. Fix: make the fixture yield and drop the table and schema on teardown. Alternatively, exclude ^delta_schema$ in the shared ingestion_config.
Clean up the Delta objects when the module finishes:
@pytest.fixture(scope="module")
def create_delta_table(trino_container):
engine = create_engine(make_url(trino_container.get_connection_url()).set(database="delta"))
try:
with engine.connect() as conn:
... # existing CREATE/INSERT statements
conn.commit()
yield
finally:
with engine.connect() as conn:
conn.execute(text("DROP TABLE IF EXISTS delta.delta_schema.delta_sales"))
conn.execute(text("DROP SCHEMA IF EXISTS delta.delta_schema"))
engine.dispose()
- Apply fix
Check the box to apply the fix or reply for a change | Was this helpful? React with 👍 / 👎
Code Review
|
| Compact |
|
Was this helpful? React with 👍 / 👎 | Powered by Gitar — free for open source
|



Describe your changes:
Part of open-metadata/openmetadata-collate#5994
I worked on giving Delta Lake tables their own table type, and on teaching the Trino and StarRocks
connectors to detect them, because today a Delta table that these connectors already see is ingested
as
ExternalorRegular— indistinguishable from a plain table, the same gap Iceberg had before itgot its own type.
The schema change lands first and on its own: the server rejects an unknown
tableTypewith400 Invalid request formatat the sink, so no connector change is testable end to end until theenum exists.
Detection uses only what each connector's own catalog already reports. No connector reads
_delta_log, and neither connector issues an extra query per table.Type of change:
High-level design:
Schema first.
DeltaLakeis added totableTypeinopenmetadata-spec/.../entity/data/table.json— both theenumand thejavaEnumsblock, directlyafter
Iceberg. Python models and the three committed TypeScript files that carry theTableTypeenum were regenerated with the repo's own generators (
make generate,json2ts-generate-all.sh);nothing generated was hand-edited. Running the TS generator on an unmodified tree first produced a
zero-file diff, which is the proof that the three changed files are the generator's output and not
mine.
Naming:
DeltaLake, notDelta. The service type is alreadyDeltaLake, the connector directory isdeltalake, and the enum already holds compound values (MaterializedView,SecureView). BareDeltareads as a diff in a metadata tool.Detection is connector-owned. No shared helper was introduced; each connector uses the indicator
its own catalog exposes, in the seam that already decides table type:
system.metadata.catalogs.connector_name = 'delta_lake'query_table_names_and_types, oneelifafter the existing Iceberg branchINFORMATION_SCHEMA.tables.ENGINE = 'DELTALAKE'RELKIND_MAPThe StarRocks value is upper case, measured on a live 3.2.16 instance as
HEX(ENGINE) = 44454C54414C414B45, 9 bytes. Engines with no map entry keep falling back toRegular, so no other key is touched. The Trino comparison is exact ondelta_lake:deltaisPrestoDB's connector name and must not match.
Existing entities are preserved.
TableRepositorydiffstableTypeon update, so a tablealready ingested as
RegularorExternalflips type on the same entity — no new id, no new FQN, nomigration.
TableResourceIT#patch_tableTypeRegularToDeltaLake_keepsSameEntityasserts exactly that,and both live runs confirmed it on real entities (version
0.1→0.2, id unchanged).Alternatives rejected. A shared
CATALOG_CONNECTOR_TABLE_TYPEmap for Trino was considered andskipped: two branches do not justify a registry, and converting it would rewrite the shipped Iceberg
branch. Reading
_delta_logfrom these connectors was rejected outright — it is the dedicateddeltalakeconnector's job, and it would need objectstorage credentials these connectors do not have.
Backward compatibility. Widening an enum only. Existing rows are untouched, no stored value
changes, and no migration is required.
Tests:
Use cases covered
delta_lakecatalog is ingested withtableType: DeltaLake.deltalakecatalog is ingested withtableType: DeltaLake.hivecatalog staysRegular; a StarRocksICEBERG/HIVE/TABLEenginekeeps its existing type.
delta(PrestoDB's name) does not match.Regularkeeps its id and FQN when it is re-ingested asDeltaLake.Regularinstead of raising.Unit tests
ingestion/tests/unit/topology/database/test_trino_metadata.py— 3 cases added toTestTrinoIcebergDetection(delta_lake catalog, hive catalog negative,deltamust not match).ingestion/tests/unit/topology/database/test_starrocks.py—TestStarRocksDeltaLakeDetection,a parametrized test driving the real
StarRocksSource.query_table_names_and_types(
DELTALAKE→DeltaLake, mixed case→Regular,ICEBERG/HIVE/TABLEunchanged).pytest --cov --cov-report=jsonand intersected with thediff:
trino/metadata.py— changed lines 373-374, 100% covered.starrocks/metadata.py— changed line 59, 100% covered.test_trino_metadata.py19 passed ·test_starrocks.py29 passed. Whole affected set:570 passed (every unit test touching Trino) and 82 passed (every unit test touching StarRocks).
Backend integration tests
openmetadata-integration-tests/.openmetadata-integration-tests/.../it/tests/TableResourceIT.java—patch_tableTypeRegularToDeltaLake_keepsSameEntity.Tests run: 1, Failures: 0, Errors: 0, Skipped: 0·BUILD SUCCESS.Regular, patches it toDeltaLake, and asserts the same id, the same FQNand the new type on readback.
Ingestion integration tests
ingestion/tests/integration/trino/test_delta_lake.pyingestion/tests/integration/trino/trino/etc/catalog/delta.propertiesdelta_lakecatalog over the Hive metastore the existing Trino harness already runs,writes a real Delta table through Trino, runs a real
MetadataWorkflow, and assertstableTypeand the column list from the API.
1 passed in 103.43s. With the detection branch reverted the same test reportsRegular,so it fails on behaviour rather than on the fixture.
(it reads them only), so there is nothing to extend for it. Its detection is covered by the unit
test that drives the real method, and by the live run below.
Playwright (UI) tests
src/generated/**type files and two connector help markdown files).Manual testing performed
Both connectors were verified against real engines with a real
metadata ingestand an API readback,each with a control run first (connector file reverted to
main) so the type change is attributableto this diff and not to the environment.
system.metadata.catalogsreportsdelta → delta_lakeandminio → hive._delta_logpresent in object storage) and a plain Hive table.metadata ingestwithmain's connector → readback…delta.delta_schema.delta_sales tableType=Regular cols=3.metadata ingeston the same service → readback…delta.delta_schema.delta_sales tableType=DeltaLake cols=3, same entity id, version0.1→0.2. Negative case in the same stack:…minio.hive_schema.hive_orders tableType=Regular cols=2.deltalakecatalog:control
tableType=Regular, thentableType=DeltaLake cols=3with this change, same id,version
0.1→0.2.SELECT ENGINE, HEX(ENGINE), LENGTH(ENGINE)→DELTALAKE / 44454C54414C414B45 / 9.Gates:
make py_format_check→All checks passed!,2915 files already formatted;mvn spotless:check -pl :openmetadata-integration-tests→BUILD SUCCESS;mvn test-compile -pl :openmetadata-spec,:openmetadata-service→BUILD SUCCESS;yarn parse-schema→ no drift.UI screen recording / screenshots:
Not applicable. No UI source changes; the UI files in this PR are generated type definitions and two
connector help documents.
Known limitations (documented, not introduced by this PR)
connector_nameand StarRocks'ENGINEdescribe the catalog a table is read through, not the table's storage format. On a metastore
shared by a
hiveand adelta_lakecatalog, a plain Hive table seen through the Delta catalog istyped
DeltaLakewith zero columns, because the engine then refuses the read. The shipped Icebergrule already behaves this way; Delta inherits it rather than introducing it, and both connector
docs now say so.
Partitioned.common_db_source.pyoverwrites the detectedtype for any partitioned non-view table. Pre-existing for Iceberg, untouched here, documented in
both connector docs.
connectionArgumentsinit_command: SET CATALOG <catalog>, and Database Schema must be left empty, since it is sentas the connection's database. A first-class
catalogfield is an enhancement, not a prerequisite.Checklist:
Fixes <issue-number>: <short explanation>— N/A: the issue lives inopenmetadata-collateand this PR is one slice of it, so it links with "Part of" rather than aclosing keyword.
Fixes #<issue-number>above — linked asPart of open-metadata/openmetadata-collate#5994; deliberately not a closing keyword, sincethis PR ships one slice of that issue.
not needed: this widens an enum, no stored value changes, and
TableRepositoryalready diffstableTypeon update, so existing entities flip type in place on the next ingestion.changes.
New feature:
building it.
table types each emits and the limitations above.