diff --git a/CHANGELOG.md b/CHANGELOG.md index 739c6a0e..f5ee8f67 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -77,6 +77,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 default namespace. An override that sets `iceberg.tables` or `table-namespace` itself is not affected. +### Added + +- **`catalog: lakekeeper` is checked as Lakekeeper.** A Lakekeeper, Polaris or Unity + catalog addresses its warehouse by NAME (`analytics`, or `/` on + Lakekeeper), so `fluid validate` refuses an object-store URI there, and warns when a + Lakekeeper `uri` does not end in `/catalog`, where Lakekeeper serves its Iceberg REST + API. Schema 0.7.6 (preview) accepts `lakekeeper` and `polaris` in `sink.catalog`; + 0.7.5 is unchanged. + ### Changed - CI: the emulated-heavy lane runs a real Kafka Connect worker (cp-kafka-connect 7.9.10 diff --git a/fluid_build/build_runners/kafka_connect/iceberg_sink_validation.py b/fluid_build/build_runners/kafka_connect/iceberg_sink_validation.py index 91455200..65be39ec 100644 --- a/fluid_build/build_runners/kafka_connect/iceberg_sink_validation.py +++ b/fluid_build/build_runners/kafka_connect/iceberg_sink_validation.py @@ -50,6 +50,9 @@ * a bucket ``{{ env.* }}`` template the DynamoDB, JDBC or BigQuery warehouse derives from, naming a variable that is unset or empty, is a warning here and an error in the runner's preflight; +* a catalog that addresses its warehouse by NAME (Lakekeeper, Polaris, Unity) + must not be given an object-store URI, and one that mounts its REST API + under a path (Lakekeeper's ``/catalog``) should have it in ``uri``; * the runtime must ship the catalog's client (the stock Apache Iceberg Kafka Connect runtime has no Nessie client, and the published sink predates the ``bigquery`` catalog type); @@ -106,6 +109,7 @@ catalog_kind_info, iceberg_catalog_kind, iceberg_sink_exposes, + is_object_store_uri, known_catalog_kinds, resolve_iceberg_catalog, resolve_iceberg_sink_exposes, @@ -552,6 +556,28 @@ def _check_catalog( "the connector needs iceberg.catalog.client.region" ) + # A name-addressed catalog owns the storage location: handed an s3:// URI + # as ``warehouse`` it looks up a warehouse literally called that and the + # REST client's /v1/config fails before the first commit. HARD. + warehouse = loc.get("warehouse") + if info.warehouse_is_name and is_object_store_uri(warehouse): + errors.append( + f"iceberg sink (build {bid!r}): {kind} addresses a warehouse by NAME " + '(e.g. "analytics" or "/"), not an object-store URI; ' + f"binding.location.warehouse is {warehouse!r}" + ) + + # Lakekeeper mounts the Iceberg REST API under ``/catalog``; a bare host + # sends the client's GET /v1/config to a path that does not serve it. + # Advisory: a reverse proxy may legitimately rewrite the path. + uri = str(loc.get("uri") or "") + if info.uri_suffix and uri and not uri.rstrip("/").endswith(info.uri_suffix): + warnings.append( + f"iceberg sink (build {bid!r}): binding.location.uri {uri!r} does not end in " + f"{info.uri_suffix!r}; {kind} serves the Iceberg REST API under " + f"{info.uri_suffix}" + ) + # The stock Apache Iceberg Kafka Connect runtime bundles the AWS, GCP and # Azure modules (Hive only in its -hive- distribution) but no iceberg-nessie # (apache/iceberg kafka-connect/build.gradle), so ``type=nessie`` fails to @@ -762,7 +788,7 @@ def _auto_creates(streaming: Mapping[str, Any], runtime: IcebergSinkPlan) -> boo def _requires_hint(key: str, kind: str, info: CatalogKind) -> str: """What a missing ``sink_requires`` key is to the sink, for check 4's error.""" - if key == "warehouse" and info.family == FAMILY_REST: + if key == "warehouse" and (info.warehouse_is_name or info.family == FAMILY_REST): return " (the catalog name)" if key == "warehouse" and kind in BUCKET_WAREHOUSE_KINDS: return ( diff --git a/fluid_build/providers/_iceberg_catalog.py b/fluid_build/providers/_iceberg_catalog.py index 003923de..53b82e7d 100644 --- a/fluid_build/providers/_iceberg_catalog.py +++ b/fluid_build/providers/_iceberg_catalog.py @@ -99,6 +99,11 @@ class CatalogKind: snowflake_catalog_type: Optional[str] #: ``binding.location`` keys a streaming sink needs for this catalog. sink_requires: Tuple[str, ...] = () + #: The warehouse is a warehouse/catalog NAME, never an object-store URI. + warehouse_is_name: bool = False + #: Path the catalog serves its REST API under (Lakekeeper mounts + #: ``/catalog``), or ``None`` when the catalog prescribes none. + uri_suffix: Optional[str] = None @property def speaks_rest(self) -> bool: @@ -121,6 +126,8 @@ def speaks_rest(self) -> bool: None, "iceberg_rest", _REST_REQUIRES, + warehouse_is_name=True, + uri_suffix="/catalog", ), CatalogKind( "polaris", @@ -129,6 +136,7 @@ def speaks_rest(self) -> bool: None, "iceberg_rest", _REST_REQUIRES, + warehouse_is_name=True, ), CatalogKind( "unity", @@ -137,6 +145,7 @@ def speaks_rest(self) -> bool: None, "iceberg_rest", _REST_REQUIRES, + warehouse_is_name=True, ), # Native NessieCatalog for the sinks (uri ends /api/v1|v2, the # warehouse is an object-store location), while Snowflake reaches diff --git a/fluid_build/schemas/fluid-schema-0.7.6.json b/fluid_build/schemas/fluid-schema-0.7.6.json index dfba70d0..775d2dd1 100644 --- a/fluid_build/schemas/fluid-schema-0.7.6.json +++ b/fluid_build/schemas/fluid-schema-0.7.6.json @@ -1943,7 +1943,7 @@ }, "catalog": { "type": "string", - "description": "NEW in v0.7.5: The Iceberg catalog that owns this table: glue, rest, lakekeeper, polaris, unity, nessie, bigquery, hive, jdbc, hadoop, dynamodb or snowflake-managed. Spellings fold case and '-'/'_', so iceberg_rest == rest and snowflake == snowflake-managed; fluid validate refuses any other value on an Iceberg expose. When absent: glue on platform aws, snowflake-managed on snowflake, rest elsewhere. The REST catalogs (rest, lakekeeper, polaris, unity) need location.uri and location.warehouse. On platform aws, a kind other than glue gets no Glue table and no Glue grant: that catalog creates and governs the table. Every emitter (streaming sinks, dbt catalogs.yml, the IaC, policy compile) reads one table, providers/_iceberg_catalog.py." + "description": "NEW in v0.7.5: The Iceberg catalog that owns this table: glue, rest, lakekeeper, polaris, unity, nessie, bigquery, hive, jdbc, hadoop, dynamodb or snowflake-managed. Spellings fold case and '-'/'_', so iceberg_rest == rest and snowflake == snowflake-managed; fluid validate refuses any other value on an Iceberg expose. When absent: glue on platform aws, snowflake-managed on snowflake, rest elsewhere. The REST catalogs (rest, lakekeeper, polaris, unity) need location.uri plus location.warehouse as the catalog's warehouse NAME, not an object-store URI; Lakekeeper serves its Iceberg REST API under /catalog (uri: http://lakekeeper:8181/catalog). On platform aws, a kind other than glue gets no Glue table and no Glue grant: that catalog creates and governs the table. Every emitter (streaming sinks, dbt catalogs.yml, the IaC, policy compile) reads one table, providers/_iceberg_catalog.py." }, "warehouse": { "type": "string", @@ -3591,7 +3591,9 @@ "nessie", "unity", "snowflake-managed", - "hive" + "hive", + "lakekeeper", + "polaris" ] }, "partitionBy": { diff --git a/tests/build_runners/test_iceberg_sink_deriver.py b/tests/build_runners/test_iceberg_sink_deriver.py index 3a1c6ee6..8293114a 100644 --- a/tests/build_runners/test_iceberg_sink_deriver.py +++ b/tests/build_runners/test_iceberg_sink_deriver.py @@ -360,21 +360,17 @@ def test_runner_preflight_warning_logs_and_still_deploys(kafka_connect_mock, tmp contract = _iceberg_contract() contract["exposes"][0]["binding"] = { **_LAKEKEEPER_NO_URI, - "location": { - **_LAKEKEEPER_NO_URI["location"], - "catalog": "nessie", - "uri": "http://nessie:19120/api/v2", - "warehouse": "s3://lake/warehouse", - }, + "location": {**_LAKEKEEPER_NO_URI["location"], "uri": "http://lakekeeper:8181"}, } with caplog.at_level("WARNING", logger="fluid.acquire.kafka_connect"): rc = execute_kafka_connect_build(contract["builds"][0], contract, tmp_path) assert rc == 0 - assert any("iceberg-nessie" in r.getMessage() for r in caplog.records) + assert any("/catalog" in r.getMessage() for r in caplog.records) cfg = kafka_connect_mock.connectors["src-sink"]["config"] - assert cfg["iceberg.catalog.type"] == "nessie" + assert cfg["iceberg.catalog.type"] == "rest" assert "iceberg.catalog.catalog-impl" not in cfg - assert cfg["iceberg.catalog.uri"] == "http://nessie:19120/api/v2" + assert cfg["iceberg.catalog.uri"] == "http://lakekeeper:8181" + assert cfg["iceberg.catalog.warehouse"] == "analytics" def test_complete_lakekeeper_streams_over_rest(kafka_connect_mock, tmp_path): diff --git a/tests/build_runners/test_iceberg_sink_validation.py b/tests/build_runners/test_iceberg_sink_validation.py index 46d97bb3..203c588c 100644 --- a/tests/build_runners/test_iceberg_sink_validation.py +++ b/tests/build_runners/test_iceberg_sink_validation.py @@ -226,6 +226,7 @@ def test_confluent_expose_not_treated_as_kc_sink_target(): CANONICAL_KINDS = sorted({canonical_catalog_kind(k) for k in known_catalog_kinds()}) REQUIRING_KINDS = [k for k in CANONICAL_KINDS if catalog_kind_info(k).sink_requires] +NAME_ADDRESSED_KINDS = [k for k in CANONICAL_KINDS if catalog_kind_info(k).warehouse_is_name] COMPLETE = {"uri": "http://lakekeeper:8181/catalog", "warehouse": "analytics"} @@ -370,12 +371,54 @@ def test_unknown_kind_errors_naming_value_and_listing_kinds(where): assert not any("disagrees" in e for e in errs) +# ── name-addressed warehouses and the Lakekeeper /catalog mount ───────────── + + +@pytest.mark.parametrize("kind", NAME_ADDRESSED_KINDS) +@pytest.mark.parametrize( + "warehouse", ["s3://lake/wh", "gs://lake/wh", "abfss://c@a.dfs.core.windows.net/w"] +) +def test_name_addressed_kind_rejects_object_store_warehouse(kind, warehouse): + binding = _loc_binding(kind, uri="http://c:8181/catalog", warehouse=warehouse) + errs = _errs(_sink_contract(binding)) + assert any(f"{kind} addresses a warehouse by NAME" in e and warehouse in e for e in errs), errs + + +@pytest.mark.parametrize("warehouse", ["analytics", "0190a7c2-project/analytics"]) +def test_name_addressed_kind_accepts_a_name(warehouse): + binding = _loc_binding("lakekeeper", uri="http://c:8181/catalog", warehouse=warehouse) + assert not any("by NAME" in e for e in _errs(_sink_contract(binding))) + + def test_generic_rest_accepts_an_object_store_warehouse(): # Plain REST catalogs (the apache/iceberg-rest-fixture) take an s3:// warehouse. binding = _loc_binding("rest", uri="http://iceberg:8181", warehouse="s3://bucket/warehouse/") assert validate_iceberg_sink(_sink_contract(binding)) == ([], []) +@pytest.mark.parametrize( + "uri, warns", + [ + ("http://lakekeeper:8181", True), + ("http://lakekeeper:8181/", True), + ("http://lakekeeper:8181/api", True), + ("http://lakekeeper:8181/catalog", False), + ("http://lakekeeper:8181/catalog/", False), + ], +) +def test_lakekeeper_uri_without_catalog_mount_warns(uri, warns): + binding = _loc_binding("lakekeeper", uri=uri, warehouse="analytics") + errors, warnings = validate_iceberg_sink(_sink_contract(binding)) + assert errors == [] + hit = [w for w in warnings if "serves the Iceberg REST API under /catalog" in w] + assert bool(hit) == warns, warnings + + +def test_uri_suffix_warning_is_lakekeeper_only(): + binding = _loc_binding("polaris", uri="http://polaris:8181/api", warehouse="analytics") + assert validate_iceberg_sink(_sink_contract(binding)) == ([], []) + + # ── runtime support: the stock KC runtime has no Nessie client ────────────── @@ -593,10 +636,10 @@ def test_preflight_selects_only_the_named_build(): def test_preflight_clean_build_returns_none_and_logs_warnings(caplog): - binding = _loc_binding("nessie", uri="http://nessie:19120/api/v2", warehouse="s3://b/w") + binding = _loc_binding("lakekeeper", uri="http://lakekeeper:8181", warehouse="analytics") with caplog.at_level("WARNING", logger="fluid.acquire.iceberg_sink"): assert iceberg_sink_preflight(_sink_contract(binding), "ingest") is None - assert any("iceberg-nessie" in r.getMessage() for r in caplog.records) + assert any("/catalog" in r.getMessage() for r in caplog.records) def test_preflight_ignores_unknown_build_ids(): diff --git a/tests/providers/test_iceberg_catalog_kinds.py b/tests/providers/test_iceberg_catalog_kinds.py index ef0297a0..68daaee7 100644 --- a/tests/providers/test_iceberg_catalog_kinds.py +++ b/tests/providers/test_iceberg_catalog_kinds.py @@ -138,6 +138,17 @@ def test_rest_catalogs_need_uri_and_warehouse(self, kind): assert info.speaks_rest assert set(info.sink_requires) == {"uri", "warehouse"} + @pytest.mark.parametrize("kind", ["lakekeeper", "polaris", "unity"]) + def test_vendor_catalogs_address_a_warehouse_by_name(self, kind): + assert catalog_kind_info(kind).warehouse_is_name + + def test_generic_rest_accepts_a_uri_warehouse(self): + # iceberg-rest-fixture and Tabular take an s3:// warehouse. + assert not catalog_kind_info("rest").warehouse_is_name + + def test_lakekeeper_serves_under_catalog(self): + assert catalog_kind_info("lakekeeper").uri_suffix == "/catalog" + class TestKindPrecedence: @pytest.mark.parametrize( diff --git a/tests/test_iceberg_catalog_kind_schema.py b/tests/test_iceberg_catalog_kind_schema.py index 02b06675..a25203ae 100644 --- a/tests/test_iceberg_catalog_kind_schema.py +++ b/tests/test_iceberg_catalog_kind_schema.py @@ -12,12 +12,14 @@ # See the License for the specific language governing permissions and # limitations under the License. -"""Schema surface for the Iceberg catalog kinds. - -``location.catalog`` is a free string the kind table in -``providers/_iceberg_catalog.py`` classifies, and ``sink.catalog`` is an enum. -These tests pin the schema's enum and prose to the table rather than to a -hand-kept list, so a kind the table does not know cannot pass the schema. +"""Schema surface for the Iceberg catalog kinds (Lakekeeper, Polaris). + +``sink.catalog`` is an enum, and before the open 0.7.6 preview it had no +``lakekeeper`` or ``polaris`` member, so a build could name the catalog only +on the expose's free-string ``location.catalog``. 0.7.6 accepts both; 0.7.5 +is GA and frozen, so it keeps refusing them. The kind table in +``providers/_iceberg_catalog.py`` is the source of truth, so these tests also +pin the schema's enum and prose to it rather than to a hand-kept list. """ from __future__ import annotations @@ -110,6 +112,23 @@ def test_base_contract_is_valid(version: str) -> None: assert result.is_valid, result.errors +@pytest.mark.parametrize("catalog", ["lakekeeper", "polaris"]) +def test_preview_accepts_the_new_sink_catalogs(catalog: str) -> None: + contract = _contract("0.7.6", sink_catalog=catalog, location=dict(_LAKEKEEPER_LOCATION)) + contract["exposes"][0]["binding"]["location"]["catalog"] = catalog + result = _validate(contract, "0.7.6") + assert result.is_valid, f"0.7.6 refused sink.catalog {catalog!r}: {result.errors}" + + +@pytest.mark.parametrize("catalog", ["lakekeeper", "polaris"]) +def test_ga_075_still_refuses_them(catalog: str) -> None: + """0.7.5 is released and frozen: widening its enum would change what an + already-published schema accepts.""" + result = _validate(_contract("0.7.5", sink_catalog=catalog), "0.7.5") + assert not result.is_valid + assert any("sink.catalog" in e and catalog in e for e in result.errors), result.errors + + @pytest.mark.parametrize("version", ["0.7.5", "0.7.6"]) def test_location_catalog_lakekeeper_is_a_free_string(version: str) -> None: contract = _contract(version, location=dict(_LAKEKEEPER_LOCATION)) @@ -135,3 +154,4 @@ def test_location_catalog_description_lists_every_kind() -> None: canonical = [k for k in known_catalog_kinds() if canonical_catalog_kind(k) == k] missing = [k for k in canonical if k not in description] assert missing == [], f"location.catalog description omits {missing}" + assert "/catalog" in description # where Lakekeeper serves its REST API