Implementation:Risingwavelabs Risingwave DebeziumOpenLineageEmitter
| Property | Value |
|---|---|
| Component | CDC Source - Debezium OpenLineage |
| Language | Java |
| Lines | 212 |
| License | Apache 2.0 (RisingWave Labs + Debezium Authors) |
| Repository | risingwavelabs/risingwave |
Overview
DebeziumOpenLineageEmitter is a patched version of the Debezium OpenLineage integration class that serves as a facade for emitting lineage tracking events from CDC connectors within RisingWave. It provides static methods for initializing per-connector emitters and emitting lineage events at various stages of the Debezium connector lifecycle (startup, data capture, error handling, shutdown).
The class maintains a thread-safe ConcurrentHashMap of LineageEmitter instances, keyed by connector context, and uses Java's ServiceLoader mechanism to discover and instantiate emitter factory implementations at runtime. When OpenLineage integration is disabled in the configuration, the class transparently uses a NoOpLineageEmitter that performs no actual emission.
This file is placed in the io.debezium.openlineage package to override the upstream Debezium class, allowing RisingWave to include this patched version in its CDC source module.
Code Reference
Source Location
java/connector-node/risingwave-source-cdc/src/main/java/io/debezium/openlineage/DebeziumOpenLineageEmitter.java
Signature
public class DebeziumOpenLineageEmitter
Key Methods
// Initializes the emitter for a specific connector; must be called before emission
public static void init(Map<String, String> configuration, String connectorTypeName)
// Removes the emitter for a connector on shutdown
public static void cleanup(ConnectorContext connectorContext)
// Creates a ConnectorContext from configuration and connector name
public static ConnectorContext connectorContext(Map<String, String> config, String connectorName)
// Emits a lineage event for a given task state
public static void emit(ConnectorContext connectorContext, DebeziumTaskState state)
// Emits a lineage event with an exception (error reporting)
public static void emit(ConnectorContext connectorContext, DebeziumTaskState state, Throwable t)
// Emits a lineage event with dataset metadata
public static void emit(ConnectorContext connectorContext, DebeziumTaskState state,
List<DatasetMetadata> datasetMetadata)
// Emits a lineage event with dataset metadata and an exception
public static void emit(ConnectorContext connectorContext, DebeziumTaskState state,
List<DatasetMetadata> datasetMetadata, Throwable t)
Internal Methods
// Retrieves the emitter for a connector, falling back to NoOpLineageEmitter if disabled
private static LineageEmitter getEmitter(ConnectorContext connectorContext)
// Checks whether OpenLineage integration is disabled in the configuration
private static boolean isOpenLineageDisabled(Map<String, String> configuration)
Imports
import static io.debezium.openlineage.OpenLineageConfig.OPEN_LINEAGE_INTEGRATION_ENABLED;
import io.debezium.connector.common.DebeziumTaskState;
import io.debezium.openlineage.dataset.DatasetMetadata;
import io.debezium.openlineage.emitter.LineageEmitter;
import io.debezium.openlineage.emitter.LineageEmitterFactory;
import io.debezium.openlineage.emitter.NoOpLineageEmitter;
import java.util.List;
import java.util.Map;
import java.util.ServiceLoader;
import java.util.concurrent.ConcurrentHashMap;
I/O Contract
Input
- configuration (
Map<String, String>): Debezium connector configuration containing theOPEN_LINEAGE_INTEGRATION_ENABLEDflag and other OpenLineage settings. - connectorTypeName (String): The name of the connector type (e.g.,
"mysql","postgresql"). - ConnectorContext: Encapsulates connector configuration and identity for emitter lookup.
- DebeziumTaskState: The current lifecycle state of the Debezium source task.
- DatasetMetadata (List): Metadata about input datasets being captured.
- Throwable: An optional exception for error lineage events.
Output
- Lineage events: Emitted through the underlying
LineageEmitterimplementation (or silently dropped byNoOpLineageEmitter).
Side Effects
- ConcurrentHashMap management: Stores and removes emitter instances per connector context.
- ServiceLoader discovery: Uses
ServiceLoaderwith synchronization to discoverLineageEmitterFactoryimplementations. - IllegalStateException: Thrown if
emit()is called beforeinit()for a given connector.
Usage Examples
Initializing and Emitting Lineage Events
// Initialize the emitter during connector startup
Map<String, String> config = getConnectorConfig();
DebeziumOpenLineageEmitter.init(config, "mysql");
// Create a connector context for subsequent emissions
ConnectorContext ctx = DebeziumOpenLineageEmitter.connectorContext(config, "mysql");
// Emit a lineage event for a task state
DebeziumOpenLineageEmitter.emit(ctx, DebeziumTaskState.RUNNING);
// Emit a lineage event with dataset metadata
List<DatasetMetadata> datasets = collectDatasetMetadata();
DebeziumOpenLineageEmitter.emit(ctx, DebeziumTaskState.RUNNING, datasets);
// Emit an error lineage event
DebeziumOpenLineageEmitter.emit(ctx, DebeziumTaskState.FAILED, exception);
// Clean up on shutdown
DebeziumOpenLineageEmitter.cleanup(ctx);
Related Pages
- DbzCdcEngineRunner Start - CDC engine runner that integrates with the OpenLineage emitter
- ConfigurableOffsetBackingStore - Offset backing store in the same CDC pipeline
- JniDbzSourceHandler RunJniDbzSourceThread - JNI bridge for Debezium source threads