Skip to content
Merged
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
130 changes: 130 additions & 0 deletions .github/workflows/integration-emulated-heavy.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,9 @@
# pubsub, via tests/iac/_gcp_emulator/docker-compose.yml. Keyless.
# • AWS — LocalStack (S3 + Glue). Needs LOCALSTACK_AUTH_TOKEN; Glue is a
# LocalStack Pro feature (Ultimate tier and above).
# • Lakekeeper + Kafka Connect — forge-derived Iceberg sink configs on a
# real worker against a real Lakekeeper REST catalog. Keyless, in its
# own job (`lakekeeper-integration`).
#
# WHY THIS LANE RUNS NIGHTLY, NOT ON A LABEL AND NOT PER MERGE
# ------------------------------------------------------------
Expand Down Expand Up @@ -310,3 +313,130 @@ jobs:
if: always()
run: |
datahub docker nuke --keep-data 2>&1 | tail -5 || true

# ────────────────────────────────────────────────────────────────────
# Lakekeeper + Kafka Connect lane — forge's derived Iceberg sink
# configs on a real worker.
#
# What tests/integration/test_lakekeeper_kafka_connect_live.py proves:
# - a `location.catalog: lakekeeper` expose, derived through forge's
# real code, streams records into a real Lakekeeper REST catalog,
# which vends STS credentials for Silo (S3 + STS);
# - the same worker refuses a sink carrying both
# `iceberg.catalog.type` and `iceberg.catalog.catalog-impl`, the
# startup crash every Glue sink hit before the deriver emitted
# impl XOR type;
# - forge's derived Glue config gets past that refusal.
#
# Keyless: the stack is self-contained (Postgres, Lakekeeper, Silo,
# KRaft Kafka, cp-kafka-connect, every image pinned by digest), so
# there is no `environment:` and no deployment record. It still gates
# on the same label as its siblings so un-vouched PR code gets no
# runner.
#
# The test downloads the Apache Iceberg sink ZIP itself and verifies
# its sha1. The unzipped plugin is cached here by version + sha1; if
# the test's pin moves and this key does not, the stale cache is
# simply re-downloaded over (its .sha1 marker no longer matches).
lakekeeper-integration:
name: Lakekeeper + Kafka Connect integration
runs-on: ubuntu-latest
timeout-minutes: 45
if: >-
github.event_name == 'schedule' ||
github.event_name == 'workflow_dispatch' ||
contains(github.event.pull_request.labels.*.name, 'ci:integration-emulated')
permissions:
contents: read
env:
ICEBERG_SINK_VERSION: "1.9.2"
# The ZIP's published sha1, not a credential.
ICEBERG_SINK_SHA1: 52dab2ff2b9659008deec803d1c1e92813c5ffa7 # pragma: allowlist secret
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0
with:
python-version: "3.12"
cache: pip

# The test imports only forge's pure derivation modules and drives the
# stack through `docker compose` and urllib: the base install plus the
# pytest plugin pyproject.toml configures is all it needs.
- name: Install fluid-build
run: |
python -m pip install --upgrade pip
python -m pip install -e . "pytest>=7.4" "pytest-timeout>=2.3"

- name: Restore the Iceberg sink plugin
id: plugin-cache
uses: actions/cache/restore@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 # v6.1.0
with:
path: ${{ runner.temp }}/iceberg-kafka-connect
key: iceberg-kafka-connect-${{ env.ICEBERG_SINK_VERSION }}-${{ env.ICEBERG_SINK_SHA1 }}

- name: Run the Lakekeeper + Kafka Connect live tests
timeout-minutes: 35
env:
FLUID_TEST_LAKEKEEPER: "1"
FLUID_LK_PLUGIN_CACHE: ${{ runner.temp }}/iceberg-kafka-connect
# A fixed project name so the steps below can collect logs from, and
# tear down, a stack the test could not clean up (a timeout kill).
FLUID_LK_PROJECT: fluid-lk-ci
FLUID_LK_LOG_DIR: ${{ runner.temp }}/lakekeeper-logs
run: |
python -m pytest -v -m emulated_heavy \
tests/integration/test_lakekeeper_kafka_connect_live.py \
--junitxml=lakekeeper.xml

# pytest exits 0 on an all-skipped run. This does not.
- name: Assert the lane actually exercised Lakekeeper
if: always()
run: |
python scripts/ci/assert_lane_coverage.py lakekeeper.xml \
--require "Lakekeeper + Kafka Connect=tests/integration/test_lakekeeper_kafka_connect_live.py"

# Saved whether or not the tests passed: the download does not depend on
# them. Only a verified install is saved; the test's one-shot writes the
# .sha1 marker last, after `sha1sum -c` passed.
- name: Check the plugin cache is complete
id: plugin-ready
if: always() && steps.plugin-cache.outputs.cache-hit != 'true'
env:
CACHE_DIR: ${{ runner.temp }}/iceberg-kafka-connect
run: |
if [ "$(cat "$CACHE_DIR/.sha1" 2>/dev/null)" = "$ICEBERG_SINK_SHA1" ]; then
echo "ready=true" >> "$GITHUB_OUTPUT"
fi

- name: Save the Iceberg sink plugin
if: always() && steps.plugin-ready.outputs.ready == 'true'
uses: actions/cache/save@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 # v6.1.0
with:
path: ${{ runner.temp }}/iceberg-kafka-connect
key: ${{ steps.plugin-cache.outputs.cache-primary-key }}

# The test writes compose logs to FLUID_LK_LOG_DIR when it fails; this
# catches the case it could not, a run killed before its teardown.
- name: Collect logs from a stack left behind
if: failure()
env:
LOG_DIR: ${{ runner.temp }}/lakekeeper-logs
run: |
mkdir -p "$LOG_DIR"
docker compose -p fluid-lk-ci logs --no-color --tail=500 \
> "$LOG_DIR/fluid-lk-ci-left-behind.log" 2>&1 || true

- name: Upload JUnit + compose logs
if: failure()
uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
with:
name: lakekeeper-integration-logs
path: |
lakekeeper.xml
${{ runner.temp }}/lakekeeper-logs/
if-no-files-found: ignore
retention-days: 14

- name: Tear down a stack left behind
if: always()
run: docker compose -p fluid-lk-ci down -v --remove-orphans || true
66 changes: 66 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,26 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Upgrade notes

- An AWS contract that declared `location.catalog` other than `glue` (for example
`rest`, `iceberg_rest`, `polaris`, `unity`, `nessie` or `lakekeeper`) on an Iceberg
expose holds a Glue database and table in its OpenTofu state that earlier releases
created for it. `fluid apply` now stops before planning with
`iceberg_catalog_move_blocked` and prints the `tofu state rm` commands that release
them, instead of planning to destroy them: destroying a Glue database deletes every
table in it. The resources stay in AWS; run the commands, then apply again.

### Changed

- CI: the emulated-heavy lane runs a real Kafka Connect worker (cp-kafka-connect 7.9.10
with the Apache Iceberg sink 1.9.2) against a real Lakekeeper (v0.13.6, credentials
vended over STS from an S3-compatible store). It streams records through the config
forge derives for `catalog: lakekeeper`, shows that a config carrying both
`iceberg.catalog.type` and `catalog-impl` fails on the worker, and shows forge's Glue
config gets past that check. `assert_lane_coverage.py` fails the job if every test
skipped. Run it locally with `FLUID_TEST_LAKEKEEPER=1` (Docker required).

- CI: `release.yml` can publish through a TestPyPI outage. A manual run with
`skip_testpypi: true` skips the TestPyPI upload and its install check and
publishes the existing tag straight to PyPI; `verify-pypi` still installs and
Expand All @@ -29,6 +47,54 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Fixed

- **A Kafka Connect Iceberg sink on AWS Glue starts.** The derived connector config set
both `iceberg.catalog.type` and `iceberg.catalog.catalog-impl`, and Apache Iceberg's
`CatalogUtil` refuses that ("both type and catalog-impl are set"), so the sink failed
at startup. It now sets `catalog-impl` or `type`, never both, as the Debezium Server
sink already did. `fluid validate` also refuses an `iceberg_catalog_overrides` or
hand-written sink config that would add the other key back.
- **Every emitter reads `location.catalog` the same way.** The streaming sinks, dbt
`catalogs.yml`, the Snowflake, AWS and Confluent IaC, the native AWS planner, `fluid
policy compile`, `fluid diff` and `fluid test` each classified the free-string value by
hand, and disagreed: `catalog: lakekeeper` streamed over REST while dbt wrote a
Snowflake-managed table and the AWS IaC created a Glue table of the same name. One
table in `providers/_iceberg_catalog.py` now classifies it; spellings fold case and
`-`/`_` (`iceberg-rest` is `rest`, `snowflake` is `snowflake-managed`).
- AWS: an Iceberg expose in another catalog gets its S3 bucket but no Glue database,
table, import, Glue IAM grant or Lake Formation resource, and `fluid diff` and
`fluid test` no longer look for it in Glue. `governance.lakeFormation`, column
restrictions and row filters on such an expose are refused at `fluid validate` and
at apply, by catalog name, instead of being emitted against a Glue table that does
not exist. An absent catalog or `catalog: glue` emits exactly what it did.
- Snowflake: `lakekeeper` and the `iceberg-rest` spelling are external catalogs
(`catalog_type: iceberg_rest`, no EXTERNAL VOLUME), and `fluid validate` no longer
demands an `s3://` or `gs://` warehouse for them. `hive`, `jdbc`, `hadoop` and
`dynamodb`, which Snowflake has no catalog integration for, are a validate error and
are left out of `catalogs.yml` instead of becoming a Snowflake-managed table.
- dbt-bigquery: an expose naming a catalog other than `bigquery` is left out of
`catalogs.yml` with a warning instead of becoming a BigLake table.
- Confluent Tableflow: a `location.catalog` other than `glue` is a validate error; it
was ignored and the table was published to Glue.
- `catalog: dynamodb` reaches Iceberg as `catalog-impl` (`DynamoDbCatalog`); Iceberg
has no `dynamodb` type.
- **`fluid validate` refuses an Iceberg catalog value it does not know**, and lists the
accepted ones. The emitters used to fall back differently, so a typo split one table
across catalogs. `fluid apply` on AWS refuses it too (`unknown-iceberg-catalog`).
- **The streaming-sink checks cover every catalog.** `uri` and `warehouse` were required
only for the literal `rest`; every REST catalog (and Nessie) now needs them. A
`sink.catalog` that disagrees with the expose's catalog is refused, since dbt and the
IaC read only the expose. A warehouse override equal to the REST binding's warehouse
no longer warns that "the static Glue table may differ". The Kafka Connect runner and
the embedded Debezium Server runner run the same checks before they create anything.
- **A crash in an Iceberg or Confluent gate fails `fluid validate`.** It was printed only
with `--verbose`, and the contract passed unchecked.
- **`fluid validate` catches two `catalog: snowflake` Iceberg exposes that derive one
EXTERNAL VOLUME on different storage**, instead of `fluid apply` failing mid-emit.
- **`fluid policy compile` grants an Iceberg table where its catalog lives.** Every
Iceberg expose off GCP compiled to AWS S3 and Glue grants, whatever its platform. A
Snowflake-managed table now compiles to Snowflake grants, and a table in another
catalog gets no Glue grant and a warning to enforce access in that catalog.

- **A Lake Formation tag the contract does not define is associated, and the module
validates.** A key in a binding's `governance.lakeFormation.tags` with no
`governance.lakeFormation.tagDefinitions` entry got a `depends_on` on an
Expand Down
43 changes: 38 additions & 5 deletions RFC-streaming-extension.md
Original file line number Diff line number Diff line change
Expand Up @@ -216,13 +216,22 @@ re-key is tested for zero key-loss — §10).

| Key | Value |
|---|---|
| `iceberg.catalog.type` | `glue` |
| `iceberg.catalog.type` | ~~`glue`~~ **not emitted** (correction below) |
| `iceberg.catalog.catalog-impl` | `org.apache.iceberg.aws.glue.GlueCatalog` |
| `iceberg.catalog.io-impl` | `org.apache.iceberg.aws.s3.S3FileIO` |
| `iceberg.catalog.warehouse` | the **exact** `s3://{bucket}/{path}` `get_iceberg_warehouse` builds for the static table |
| `iceberg.catalog.client.region` | from `binding.location.region` |
| credentials | **none emitted** — DefaultCredentialsProvider chain (env / instance-profile / IRSA) |

> **CORRECTION (2026-10-01).** This table listed both `iceberg.catalog.type`
> and `iceberg.catalog.catalog-impl`, while §6.8 #1 says to reject a config with
> both. §6.8 is right: Iceberg's `CatalogUtil` refuses a catalog that sets
> `type` and `catalog-impl` together. The deriver emits **`catalog-impl` XOR
> `type`**: Glue (and DynamoDB, which has no `type`) get `catalog-impl` only,
> and every other kind gets `type` only. Both come from the catalog-kind table
> in `providers/_iceberg_catalog.py` (`catalog_kind_info`), which dbt, the IaC
> emitters and the policy compiler read too.

REST (`type=rest`, warehouse = catalog **name**, requires `s3.*` creds via
`secret_ref`) and GCP (`GCSFileIO`) are **tagged-union variants deferred to
PR7**. Each variant declares its own required-key set, validated at plan time.
Expand Down Expand Up @@ -329,6 +338,18 @@ seam (`cli/plan.py`, after action-type parse, before `plan.json` serialization):
streaming-only fallback — exactly where we'd assumed "auto-create handles
it." Validator adds: streaming-only + auto-create ⇒ require a
namespace-ensure action.
**CORRECTION (2026-10-01) to the correction above.** The observation is
kept as the record of what the spike saw, but its cause was the runtime,
not the design: the spike ran the Tabular `io.tabular` **0.6.19** runtime
(§14), not the Apache sink v1 ships against. The Apache sink creates the
namespace, and each of its parents, when it auto-creates a table, since
Iceberg **1.6.0** (apache/iceberg#10186:
`IcebergWriterFactory.createNamespaceIfNotExist`, called from
`autoCreateTable`; absent in 1.5.2). So **no ensure-namespace step is
needed** and the validator requirement above is **withdrawn**. One residue:
the sink ignores a `ForbiddenException` from `createNamespace`, so a
catalog principal without create-namespace rights still needs the
namespace created for it, or the table create fails.
- **Operator override present:** when a hand-written
`iceberg.catalog.warehouse` is set, **defer to the operator (warn, not
fail)** — preserves operator-wins.
Expand Down Expand Up @@ -521,6 +542,13 @@ by diffing the 0.6.19 README against the 1.11.0 docs — so the observed run
validates the entire key surface; only those two constants differ and both are
independently doc-confirmed for 1.11.0.

**Version caveat, corrected (2026-10-01):** an identical key surface did not
mean identical behaviour. 0.6.19 predates the namespace auto-create the Apache
sink gained in Iceberg 1.6.0, which is why correction A below was observed (see
§6.8 #6). And 0.6.19 is no longer the only published runtime: the Apache
Software Foundation now publishes the Apache sink on Confluent Hub
(`iceberg/iceberg-kafka-connect`, version 1.9.2).

**Result: PASS.** 15 rows physically committed to `default.events`; schema
auto-inferred (`amount:double, name:string, id:long, region:string`); 2 snapshots;
connector + task `RUNNING`. The exact config that worked is the same shape the
Expand All @@ -532,7 +560,7 @@ converters).

| # | Observed | RFC correction |
|---|---|---|
| A | `auto-create-enabled=true` created the table but failed with `NoSuchNamespaceException: Namespace default does not exist` until the namespace was created explicitly. | The **auto-create fallback is incomplete** — forge must ensure the **namespace/database** exists even build-only. Folded into §6.8 #6. |
| A | `auto-create-enabled=true` created the table but failed with `NoSuchNamespaceException: Namespace default does not exist` until the namespace was created explicitly. | The **auto-create fallback is incomplete** — forge must ensure the **namespace/database** exists even build-only. Folded into §6.8 #6. **Withdrawn 2026-10-01:** a 0.6.19 behaviour; the Apache sink creates the namespace and its parents on auto-create since Iceberg 1.6.0 (§6.8 #6). |
| B | `connector.state=RUNNING` while `tasks[0].state=FAILED` with a full trace. | The runner **must inspect task state + trace**, not just connector state. Confirmed §11 (now marked spike-observed). |
| C | Schemaless JSON required `value.converter=JsonConverter` + `schemas.enable=false`; a task restart re-delivered records (duplicates) with no EOS. | The deriver must **co-emit converters** (§6.2, was implicit) and EOS config is **necessary not optional** (§11). |

Expand Down Expand Up @@ -630,6 +658,11 @@ reference plugin). "Engine" becomes a binding concern:
Kafka-Connect spike found (§14 A: auto-create makes the table, not the
namespace) — a satisfying cross-topology consistency: **the writer owns the
table, forge owns the namespace.**
**CORRECTION (2026-10-01):** the Tableflow conclusion stands on its own
evidence (the IAM policy has no `glue:CreateDatabase`), but the parallel does
not: §14 A was a 0.6.19 runtime behaviour, and the Apache sink creates the
namespace itself since Iceberg 1.6.0 (§6.8 #6). The split is
Tableflow-specific, not cross-topology.
- **(b) Managed-class feasibility — ANSWERED: NO.** Confluent Cloud has **no
fully-managed Iceberg sink connector**; the OSS `org.apache.iceberg.connect`
class runs only via "Bring Your Own Connector" (custom upload). The managed
Expand Down Expand Up @@ -666,7 +699,7 @@ sweep (mirrors the existing `FLUID_IAC_LIVE_{AWS,GCP,SNOWFLAKE}` gates).
| Key | Role | forge derivation |
|---|---|---|
| `connector.class` | sink class | constant `org.apache.iceberg.connect.IcebergSinkConnector` |
| `iceberg.catalog.type` / `catalog-impl` | catalog kind | from `binding.platform` |
| `iceberg.catalog.type` / `catalog-impl` | catalog kind | from the catalog kind (`location.catalog`, else the `binding.platform` default); `catalog-impl` XOR `type` (corrected 2026-10-01, §6.3) |
| `iceberg.catalog.warehouse` | warehouse root | `get_iceberg_warehouse()` (shared with static path) |
| `iceberg.catalog.io-impl` | FileIO | per-platform (`S3FileIO` for AWS) |
| `iceberg.catalog.client.region` | region | `binding.location.region` |
Expand All @@ -676,9 +709,9 @@ sweep (mirrors the existing `FLUID_IAC_LIVE_{AWS,GCP,SNOWFLAKE}` gates).
| `iceberg.tables.default-partition-by` | partitioning | `icebergConfig.partitionSpec` |
| `iceberg.control.topic` | EOS control | `_iceberg-control-{product_id}` (spike-validated custom topic) |
| `iceberg.coordinator.transactional.prefix` | EOS | `iceberg-coord-{product_id}` |
| `iceberg.tables.auto-create-enabled` | build-only fallback | `true` for streaming-only — **but namespace must pre-exist** (§14 A) |
| `iceberg.tables.auto-create-enabled` | build-only fallback | `true` for streaming-only; ~~**but namespace must pre-exist** (§14 A)~~ the sink creates the namespace and its parents too, since Iceberg 1.6.0 (corrected 2026-10-01, §6.8 #6) |
| `key.converter` / `value.converter` | record decode | `JsonConverter` + `schemas.enable=false` (schemaless JSON; §14 C) |
| (catalog namespace) | table parent | **forge must `ensure-namespace`** — auto-create does NOT (§14 A) |
| (catalog namespace) | table parent | ~~**forge must `ensure-namespace`** — auto-create does NOT (§14 A)~~ created by the sink on auto-create since Iceberg 1.6.0 (apache/iceberg#10186); no forge step (corrected 2026-10-01, §6.8 #6) |

## Appendix B — verification log

Expand Down
Loading
Loading