kapps_triplestore_interface.kafka.kafka_manager#

class kapps_triplestore_interface.kafka.kafka_manager.KafkaManager(db: GraphDB)[source]#

Bases: object

Manage GraphDB Kafka connectors via SPARQL.

Provides helpers to list, inspect, create, and drop Kafka connectors stored in GraphDB following Ontotext’s connector ontology.

Reference:

https://graphdb.ontotext.com/documentation/11.1/kafka-graphdb-connector.html

get_existing_connector_ids() → List[str][source]#

Get the IDs of existing Kafka connectors.

Queries the graph database for resources with the kafka:listConnectors predicate.

Returns:

Connector IDs; empty when none are found.

Return type:

List[str]

Raises:

GraphDbException – If the underlying query execution fails.

get_status_of_connectors(id: str | None = None) → Dict[str, Dict] | None[source]#

Get the status of Kafka connectors.

When id is provided, returns the status of that connector; otherwise returns statuses for all connectors.

Parameters:

id (Optional[str]) – Connector id to filter by. Defaults to None.

Returns:

Mapping of connector name to status; None when no results.

Return type:

Optional[Dict[str, Dict]]

Raises:

GraphDbException – If the underlying query execution fails.

get_connector_create_options(id: str) → str | None[source]#

Retrieve the creation options for a Kafka connector.

Queries the graph database for the stored creation configuration string for the given connector instance.

Parameters:

id (str) – The connector identifier.

Returns:

The creation options string when available, otherwise None.

Return type:

Optional[str]

Raises:

GraphDbException – If the underlying query execution fails.

drop_connector(id: str) → bool[source]#

Drop the specified Kafka connector.

Parameters:

id (str) – Connector identifier to drop.

Returns:

True on success, False otherwise.

Return type:

bool

create_connector(id: str, connector_config: dict, overwrite: bool | None = False) → None[source]#

Create a Kafka connector with the specified configuration.

Inserts the connector configuration into GraphDB via an INSERT DATA query. If a connector with the same id exists and overwrite=True, it will be dropped first.

Parameters:
  • id (str) – Unique connector identifier.

  • connector_config (dict) – Kafka connector configuration to serialize and store.

  • overwrite (Optional[bool]) – Drop any existing connector with the same id first. Defaults to False.