After Kafka 129 , BookNest's lakehouse pipeline has three steps: the Connect sink lands the topic, Spark 129 builds daily_sales, and dbt 37,942 builds fct_order_lifecycle. The sink has no OpenLineage integration and the Spark listener missed its datasets, so the pipeline reports those two runs itself:
"""Report the BookNest runs that have no working OpenLineage integration on this host."""
from datetime import datetime, timezone
from openlineage.client import OpenLineageClient
from openlineage.client.event_v2 import InputDataset, Job, OutputDataset, Run, RunEvent
from openlineage.client.event_v2 import RunState
from openlineage.client.transport.http import HttpConfig, HttpTransport
from openlineage.client.uuid import generate_new_uuid
client = OpenLineageClient(transport=HttpTransport(HttpConfig(url="http://localhost:31500")))
PRODUCER = "https://example.com/booknest/pipeline" # who sends these events
def report(job, inputs, outputs):
run = Run(runId=str(generate_new_uuid()))
for state in (RunState.START, RunState.COMPLETE):
client.emit(RunEvent(eventType=state, eventTime=datetime.now(timezone.utc).isoformat(),
run=run, job=Job(namespace="booknest", name=job), inputs=inputs,
outputs=outputs, producer=PRODUCER))
print(f"{job}: {[d.name for d in inputs]} -> {[d.name for d in outputs]}")
def lake(name, cls): # Iceberg tables, named by location
return cls(namespace="s3://warehouse", name=f"booknest/{name}")
report("kafka_connect.booknest-events-sink",
[InputDataset(namespace="kafka://l1-kafka:9092", name="booknest.order-events")],
[lake("order_events", OutputDataset)])
report("spark.daily_sales", [lake("orders", InputDataset), lake("order_items", InputDataset)],
[lake("daily_sales", OutputDataset)])kafka_connect.booknest-events-sink: ['booknest.order-events'] -> ['booknest/order_events'] spark.daily_sales: ['booknest/orders', 'booknest/order_items'] -> ['booknest/daily_sales']
For dbt the real integration works: governance/dbt_ol.sh runs dbt Tests Revisited's project as dbt-ol run, which reports each model from dbt's artifacts (Emitted 4 OpenLineage events). Marquez recorded fct_order_lifecycle reading lake.booknest.order_events, but in the namespace duckdb:///home/dev/v7-l1/ch08/quality/dbt_lake.duckdb (the DuckDB 61,228 file dbt connects to), not s3://warehouse, so the two halves of the graph do not join. That is lineage's practical weak point: every producer must name a table the same way, and OpenLineage's naming conventions define no name for DuckDB or for an Iceberg 129 table as such. Agree on one name per table, add symlink facets for the aliases, and inspect the graph after every new integration.