kapps_triplestore_interface.kafka.kafka_manager#
- class kapps_triplestore_interface.kafka.kafka_manager.KafkaManager(db: GraphDB)[source]#
Bases:
objectManage GraphDB Kafka connectors via SPARQL.
Provides helpers to list, inspect, create, and drop Kafka connectors stored in GraphDB following Ontotext’s connector ontology.
- 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.
- 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.