Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions janusgraph-cdc-extension/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,13 @@
<version>5.3.1</version>
<scope>test</scope>
</dependency>

<!-- Kafka Client -->
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.4.0</version>
</dependency>
</dependencies>
<build>
<plugins>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ public class GraphLogProcessor {
private List<EventSink> sinks = new ArrayList<>();
private MessageConverter converter;
private boolean isStarted = false;
private LogProcessorFramework logProcessorFramework;

// Event buffering removed

Expand Down Expand Up @@ -97,6 +98,11 @@ private void init(JanusGraph graph, Map<String, Object> config) {
if ("LOG".equalsIgnoreCase(sinkType.trim())) {
sinks.add(new LogFileEventSink());
logger.info("Added Log File Event Sink");
} else if ("KAFKA".equalsIgnoreCase(sinkType.trim())) {
String bootstrapServers = (String) config.getOrDefault("kafka.bootstrap.servers", "kafka:9092");
String kafkaTopic = (String) config.getOrDefault("kafka.topics.graph.event", "test.knowlg.learning.graph.events");
sinks.add(new KafkaEventSink(bootstrapServers, kafkaTopic));
logger.info("Added Kafka Event Sink for topic: {}", kafkaTopic);
}
}

Expand All @@ -105,8 +111,9 @@ private void init(JanusGraph graph, Map<String, Object> config) {
}

try {
LogProcessorFramework framework = JanusGraphFactory.openTransactionLog(graph);
framework.addLogProcessor(LOG_IDENTIFIER)
// Store as instance field to prevent garbage collection
this.logProcessorFramework = JanusGraphFactory.openTransactionLog(graph);
this.logProcessorFramework.addLogProcessor(LOG_IDENTIFIER)
.setProcessorIdentifier("janusgraph-cdc-processor")
.setStartTime(Instant.now().minus(1, ChronoUnit.MINUTES))
.addProcessor(new ChangeProcessor() {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
package org.sunbird.janusgraph.cdc;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Map;
import java.util.Properties;

public class KafkaEventSink implements EventSink {

private static final Logger logger = LoggerFactory.getLogger(KafkaEventSink.class);
private KafkaProducer<String, String> producer;
private String topic;

public KafkaEventSink(String bootstrapServers, String topic) {
this.topic = topic;
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.ACKS_CONFIG, "1");
props.put(ProducerConfig.RETRIES_CONFIG, 3);
this.producer = new KafkaProducer<>(props);
logger.info("KafkaEventSink initialized with topic: {}", topic);
}

@Override
public void send(String key, String message) {
try {
producer.send(new ProducerRecord<>(topic, key, message));
} catch (Exception e) {
logger.error("Failed to send event to Kafka topic {}: {}", topic, e.getMessage());
}
}

@Override
public void close() {
if (producer != null) {
producer.flush();
producer.close();
}
}
}