From ada19c0ae05fa9bc99805c21c4ca557a6e275d6d Mon Sep 17 00:00:00 2001 From: Kartheek Palla Date: Wed, 18 Mar 2026 00:29:31 +0530 Subject: [PATCH 1/3] Filter out-of-order/duplicate graph events Add a lastUpdatedCache map and a shouldProcessEvent(...) helper to drop events that are older than or equal to the last-processed lastUpdatedOn for the same nodeUniqueId. The processor now calls shouldProcessEvent before sending events to sinks; accepted events update the cache, while dropped events are logged at debug level. This prevents duplicate or out-of-order event emissions when timestamps are present. --- .../janusgraph/cdc/GraphLogProcessor.java | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java b/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java index 5f6d543..d3afdfa 100644 --- a/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java +++ b/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java @@ -39,6 +39,10 @@ public class GraphLogProcessor { private MessageConverter converter; private boolean isStarted = false; + // Tracks the last-processed lastUpdatedOn (epoch ms) per nodeUniqueId + // to filter out-of-order/duplicate events. + private final Map lastUpdatedCache = new HashMap<>(); + // Event buffering removed private GraphLogProcessor() { @@ -219,6 +223,10 @@ private void processVertexChange(JanusGraphVertex vertex, ChangeState changeStat return; } + if (!shouldProcessEvent(event)) { + return; + } + // 4. Send event immediately (No buffering) sendEventToSinks(vertex.id().toString(), event); @@ -340,4 +348,28 @@ private Long parseTimestamp(Object ts) { } return null; } + + /** + * Returns true if this event should be emitted; false if it is older than (or + * the same age as) an already-processed event for the same vertex. + * Updates the cache when the event is accepted. + */ + private boolean shouldProcessEvent(Map event) { + String nodeId = getNodeUniqueId(event); + if (nodeId == null) { + return true; + } + Long eventTs = getLastUpdatedOn(event); + if (eventTs == null) { + return true; + } + Long cachedTs = lastUpdatedCache.get(nodeId); + if (cachedTs == null || eventTs > cachedTs) { + lastUpdatedCache.put(nodeId, eventTs); + return true; + } + logger.debug("Dropping out-of-order/duplicate event for node {} (event ts={}, cached ts={})", + nodeId, eventTs, cachedTs); + return false; + } } From 5d3e213ccda9cb6f841d9e32723b50e1ad8c7c3a Mon Sep 17 00:00:00 2001 From: Kartheek Palla Date: Mon, 23 Mar 2026 14:24:04 +0530 Subject: [PATCH 2/3] Bound cache, update on success, add /src to gitignore Make lastUpdatedCache a synchronized LRU (LinkedHashMap) with MAX_CACHE_SIZE=10_000 to prevent unbounded heap growth. Modify processing flow so events are only recorded in the cache after a confirmed successful send: sendEventToSinks now returns boolean, shouldProcessEvent no longer mutates the cache, and a new updateLastUpdatedCache helper advances the cache after success. Also clear the cache on stop. Add /src to .gitignore. These changes reduce lost retries on serialization/send failures and cap memory usage. --- .../janusgraph/cdc/GraphLogProcessor.java | 48 ++++++++++++++++--- 1 file changed, 42 insertions(+), 6 deletions(-) diff --git a/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java b/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java index d3afdfa..0cfd93a 100644 --- a/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java +++ b/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java @@ -41,7 +41,15 @@ public class GraphLogProcessor { // Tracks the last-processed lastUpdatedOn (epoch ms) per nodeUniqueId // to filter out-of-order/duplicate events. - private final Map lastUpdatedCache = new HashMap<>(); + // Bounded to MAX_CACHE_SIZE entries (LRU eviction) to prevent unbounded heap growth. + private static final int MAX_CACHE_SIZE = 10_000; + private final Map lastUpdatedCache = Collections.synchronizedMap( + new LinkedHashMap(16, 0.75f, true) { + @Override + protected boolean removeEldestEntry(Map.Entry eldest) { + return size() > MAX_CACHE_SIZE; + } + }); // Event buffering removed @@ -146,6 +154,7 @@ private void stop() { } } sinks.clear(); + lastUpdatedCache.clear(); isStarted = false; logger.info("GraphLogProcessor stopped."); } @@ -228,7 +237,12 @@ private void processVertexChange(JanusGraphVertex vertex, ChangeState changeStat } // 4. Send event immediately (No buffering) - sendEventToSinks(vertex.id().toString(), event); + // Update cache only after a confirmed successful send to avoid + // dropping retries when the send fails. + boolean sent = sendEventToSinks(vertex.id().toString(), event); + if (sent) { + updateLastUpdatedCache(event); + } } catch (Exception e) { logger.error("Error converting/processing vertex change event", e); @@ -236,9 +250,13 @@ private void processVertexChange(JanusGraphVertex vertex, ChangeState changeStat } /** - * Send event to all configured sinks + * Send event to all configured sinks. + * + * @return true if serialization succeeded (even if individual sinks erred), + * false if the event could not be serialized (so the cache should not + * be advanced and the event may be retried). */ - private void sendEventToSinks(String key, Map event) { + private boolean sendEventToSinks(String key, Map event) { try { String json = mapper.writeValueAsString(event); for (EventSink sink : sinks) { @@ -249,8 +267,10 @@ private void sendEventToSinks(String key, Map event) { } } logger.info("Sent event: {}", json); + return true; } catch (Exception e) { logger.error("Error serializing event", e); + return false; } } @@ -352,7 +372,8 @@ private Long parseTimestamp(Object ts) { /** * Returns true if this event should be emitted; false if it is older than (or * the same age as) an already-processed event for the same vertex. - * Updates the cache when the event is accepted. + * Does NOT update the cache — call {@link #updateLastUpdatedCache} after a + * confirmed successful send. */ private boolean shouldProcessEvent(Map event) { String nodeId = getNodeUniqueId(event); @@ -365,11 +386,26 @@ private boolean shouldProcessEvent(Map event) { } Long cachedTs = lastUpdatedCache.get(nodeId); if (cachedTs == null || eventTs > cachedTs) { - lastUpdatedCache.put(nodeId, eventTs); return true; } logger.debug("Dropping out-of-order/duplicate event for node {} (event ts={}, cached ts={})", nodeId, eventTs, cachedTs); return false; } + + /** + * Advances the cache for the event's nodeUniqueId to its lastUpdatedOn + * timestamp. Call this only after the event has been successfully sent. + */ + private void updateLastUpdatedCache(Map event) { + String nodeId = getNodeUniqueId(event); + if (nodeId == null) { + return; + } + Long eventTs = getLastUpdatedOn(event); + if (eventTs == null) { + return; + } + lastUpdatedCache.put(nodeId, eventTs); + } } From c77c69429e181212237e54d9b9411b931fda4a21 Mon Sep 17 00:00:00 2001 From: Kartheek Palla Date: Mon, 23 Mar 2026 14:30:49 +0530 Subject: [PATCH 3/3] fix: address PR review comments and additional code quality issues - Remove dead hasStatusAttribute method (never called) - Cache addedVertices/removedVertices sets before UPDATE loop to avoid repeated getVertices() calls and re-consuming single-use iterables - Downgrade per-edge relation logging from INFO to DEBUG to reduce noise - Fix typo: INitialize -> Initialize in comment Co-Authored-By: Claude Sonnet 4.6 --- .../janusgraph/cdc/GraphLogProcessor.java | 42 ++++++------------- 1 file changed, 13 insertions(+), 29 deletions(-) diff --git a/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java b/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java index 0cfd93a..449f69f 100644 --- a/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java +++ b/janusgraph-cdc-extension/src/main/java/org/sunbird/janusgraph/cdc/GraphLogProcessor.java @@ -87,7 +87,7 @@ private void init(JanusGraph graph, Map config) { logger.info("Starting GraphLogProcessor..."); - // INitialize Converter + // Initialize Converter String converterType = (String) config.getOrDefault("graph.txn.log_processor.converter", "DEFAULT"); if ("TELEMETRY".equalsIgnoreCase(converterType)) { converter = new TelemetryMessageConverter(); @@ -180,33 +180,37 @@ private void processChanges(TransactionId txId, ChangeState changeState) { // This ensures edge-only changes are also discovered. Set changedVertices = changeState.getVertices(Change.ANY); - // Count all relations in this ChangeState (edges + properties combined) + // Cache ADDED/REMOVED vertex sets once — avoids repeated getVertices() calls + // inside the loop and prevents re-consuming potentially single-use iterables. + Set addedVertices = changeState.getVertices(Change.ADDED); + Set removedVertices = changeState.getVertices(Change.REMOVED); + + // Count all relations in this ChangeState for diagnostic logging int addedRelCount = 0; int removedRelCount = 0; for (JanusGraphRelation r : changeState.getRelations(Change.ADDED)) { addedRelCount++; if (r.isEdge()) { - logger.info(" ADDED relation (edge): type={}, vertices={}", r.getType().name(), r); + logger.debug(" ADDED relation (edge): type={}, vertices={}", r.getType().name(), r); } } for (JanusGraphRelation r : changeState.getRelations(Change.REMOVED)) { removedRelCount++; if (r.isEdge()) { - logger.info(" REMOVED relation (edge): type={}, vertices={}", r.getType().name(), r); + logger.debug(" REMOVED relation (edge): type={}, vertices={}", r.getType().name(), r); } } logger.info("ChangeState — vertices ADDED: {}, REMOVED: {}, ANY: {} | relations ADDED: {}, REMOVED: {}", - changeState.getVertices(Change.ADDED).size(), - changeState.getVertices(Change.REMOVED).size(), - changedVertices.size(), + addedVertices.size(), removedVertices.size(), changedVertices.size(), addedRelCount, removedRelCount); + for (JanusGraphVertex vertex : changedVertices) { // If it's a new vertex, we already processed it as CREATE - if (changeState.getVertices(Change.ADDED).contains(vertex)) { + if (addedVertices.contains(vertex)) { continue; } // If it's a removed vertex, we already processed it as DELETE - if (changeState.getVertices(Change.REMOVED).contains(vertex)) { + if (removedVertices.contains(vertex)) { continue; } @@ -274,26 +278,6 @@ private boolean sendEventToSinks(String key, Map event) { } } - private boolean hasStatusAttribute(Map event) { - try { - if (event.containsKey("transactionData")) { - Map txData = (Map) event.get("transactionData"); - if (txData != null && txData.containsKey("properties")) { - Map props = (Map) txData.get("properties"); - return props != null && props.containsKey("status"); - } - } - // Check for flat properties (SimpleMessageConverter) - if (event.containsKey("properties")) { - Map props = (Map) event.get("properties"); - return props != null && props.containsKey("status"); - } - } catch (Exception e) { - // ignore - } - return false; - } - private String getNodeUniqueId(Map event) { if (event.containsKey("nodeUniqueId")) { return (String) event.get("nodeUniqueId");