Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<project-id>/<name>` 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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 "<project-id>/<name>"), 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"<host>{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
Expand Down Expand Up @@ -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 (
Expand Down
9 changes: 9 additions & 0 deletions fluid_build/providers/_iceberg_catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -121,6 +126,8 @@ def speaks_rest(self) -> bool:
None,
"iceberg_rest",
_REST_REQUIRES,
warehouse_is_name=True,
uri_suffix="/catalog",
),
CatalogKind(
"polaris",
Expand All @@ -129,6 +136,7 @@ def speaks_rest(self) -> bool:
None,
"iceberg_rest",
_REST_REQUIRES,
warehouse_is_name=True,
),
CatalogKind(
"unity",
Expand All @@ -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
Expand Down
6 changes: 4 additions & 2 deletions fluid_build/schemas/fluid-schema-0.7.6.json
Original file line number Diff line number Diff line change
Expand Up @@ -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 <host>/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",
Expand Down Expand Up @@ -3591,7 +3591,9 @@
"nessie",
"unity",
"snowflake-managed",
"hive"
"hive",
"lakekeeper",
"polaris"
]
},
"partitionBy": {
Expand Down
14 changes: 5 additions & 9 deletions tests/build_runners/test_iceberg_sink_deriver.py
Original file line number Diff line number Diff line change
Expand Up @@ -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("<host>/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):
Expand Down
47 changes: 45 additions & 2 deletions tests/build_runners/test_iceberg_sink_validation.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"}


Expand Down Expand Up @@ -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 <host>/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 ──────────────


Expand Down Expand Up @@ -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():
Expand Down
11 changes: 11 additions & 0 deletions tests/providers/test_iceberg_catalog_kinds.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
32 changes: 26 additions & 6 deletions tests/test_iceberg_catalog_kind_schema.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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))
Expand All @@ -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
Loading