Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Risingwavelabs Risingwave DebeziumOpenLineageEmitter

From Leeroopedia


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

View on GitHub

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 the OPEN_LINEAGE_INTEGRATION_ENABLED flag 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 LineageEmitter implementation (or silently dropped by NoOpLineageEmitter).

Side Effects

  • ConcurrentHashMap management: Stores and removes emitter instances per connector context.
  • ServiceLoader discovery: Uses ServiceLoader with synchronization to discover LineageEmitterFactory implementations.
  • IllegalStateException: Thrown if emit() is called before init() 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

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment