diff --git a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/TaskExecutionProcessor.java b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/TaskExecutionProcessor.java index 54d73e472..9d3290980 100644 --- a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/TaskExecutionProcessor.java +++ b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/TaskExecutionProcessor.java @@ -24,6 +24,7 @@ import org.slf4j.LoggerFactory; import java.sql.SQLException; +import java.util.List; import java.util.Objects; @Unremovable @@ -32,7 +33,7 @@ public class TaskExecutionProcessor implements EventProcessor { private static final Logger log = LoggerFactory.getLogger(TaskExecutionProcessor.class); - final TaskPersistence taskPersistence; + private final TaskPersistence taskPersistence; @Inject public TaskExecutionProcessor(TaskPersistence taskPersistence) { @@ -41,22 +42,36 @@ public TaskExecutionProcessor(TaskPersistence taskPersistence) { @Override public void process(TaskExecution event) { - Objects.requireNonNull(event, "event cannot be null"); - log.debug("Processing task: {}", event); - event.setId(generateTaskExecutionId(event)); + processBatch(List.of(event)); + } + + public void processBatch(List events) { + Objects.requireNonNull(events, "events cannot be null"); + + if (events.isEmpty()) { + return; + } + + log.debug("Processing task batch size: {}", events.size()); + try { - this.taskPersistence.persist(event); - log.debug("Successfully processed the task event with ID: {}", event.getInstanceId()); + for (TaskExecution event : events) { + Objects.requireNonNull(event, "event cannot be null"); + event.setId(generateTaskExecutionId(event)); + } + + this.taskPersistence.persistBatch(events); + + log.debug("Successfully processed task batch size: {}", events.size()); } catch (SQLException e) { - log.error("Error while processing the task event: {}", event, e); - throw new ProcessEventFailedException("Failed to process the task event with instance ID: " + event.getInstanceId(), e); + log.error("Error while processing task batch size: {}", events.size(), e); + throw new ProcessEventFailedException("Failed to process task event batch", e); } } private String generateTaskExecutionId(TaskExecution taskExecutionEvent) { - // Generate deterministic ID based on instance's ID + task position - return taskExecutionEvent.getInstanceId() + - ":" + taskExecutionEvent.getTaskPosition(); + return taskExecutionEvent.getInstanceId() + + ":" + + taskExecutionEvent.getTaskPosition(); } } - diff --git a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/WorkflowEventProcessor.java b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/WorkflowEventProcessor.java index 7063c709b..40422aee9 100644 --- a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/WorkflowEventProcessor.java +++ b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/WorkflowEventProcessor.java @@ -24,6 +24,7 @@ import org.slf4j.LoggerFactory; import java.sql.SQLException; +import java.util.List; import java.util.Objects; @Unremovable @@ -40,13 +41,19 @@ public WorkflowEventProcessor(WorkflowPersistence workflowPersistence) { } public void process(final WorkflowInstance event) { + processBatch(List.of(event)); + } + + public void processBatch(final List events) { + log.debug("Processing workflow batch size: {}", events.size()); + try { - this.workflowPersistence.persist(Objects.requireNonNull(event, "event cannot be null")); - log.debug("Successfully processed the workflow event with ID: {}", event.getId()); + workflowPersistence.persistBatch(events); + log.debug("Successfully processed {} workflow events", events.size()); } catch (SQLException e) { - log.error("Error while processing the workflow event: {}", event, e); - throw new ProcessEventFailedException("Failed to process the workflow event with instance ID: " + event.getId(), e); + log.error("Error while processing workflow event batch", e); + throw new ProcessEventFailedException("Failed to process workflow event batch", e); } - } + } } diff --git a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/TaskPersistence.java b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/TaskPersistence.java index a5bf7d6dc..d3d81a769 100644 --- a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/TaskPersistence.java +++ b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/TaskPersistence.java @@ -17,6 +17,7 @@ import java.sql.SQLException; import java.sql.Savepoint; import java.time.ZonedDateTime; +import java.util.List; import java.util.Optional; @Unremovable @@ -66,6 +67,59 @@ public void persist(TaskExecution event) throws SQLException { } } + public void persistBatch(List events) throws SQLException { + if (events == null || events.isEmpty()) { + return; + } + + log.debug("Persisting task DB batch size: {}", events.size()); + + try (Connection conn = dataSource.getConnection(); + PreparedStatement stmt = conn.prepareStatement(insertTaskUpsert)) { + + conn.setAutoCommit(false); + + try { + for (TaskExecution event : events) { + setTaskParameters(stmt, event); + stmt.addBatch(); + } + + int[] result = stmt.executeBatch(); + conn.commit(); + + log.debug("Committed task DB batch size: {}, executeBatch result length: {}", + events.size(), result.length); + } catch (SQLException e) { + conn.rollback(); + // A task may arrive before its workflow. The batch upsert cannot create the + // missing placeholder workflow, so fall back to per-record persistence which + // recovers from foreign key violations individually. + if (isForeignKeyViolation(e)) { + log.debug("Task batch hit foreign key violation; falling back to per-record persistence"); + persistEachIndividually(events); + } else { + throw e; + } + } + } + } + + private void persistEachIndividually(List events) throws SQLException { + for (TaskExecution event : events) { + persist(event); + } + } + + private boolean isForeignKeyViolation(SQLException e) { + for (SQLException current = e; current != null; current = current.getNextException()) { + if (INVALID_FOREIGN_KEY.equals(current.getSQLState())) { + return true; + } + } + return false; + } + private void tryInsertTask(TaskExecution event, Connection conn) throws SQLException { try (PreparedStatement stmt = conn.prepareStatement(insertTaskUpsert)) { setTaskParameters(stmt, event); diff --git a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/WorkflowPersistence.java b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/WorkflowPersistence.java index 6a4646e46..ccbae10a7 100644 --- a/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/WorkflowPersistence.java +++ b/data-index/data-index-ingestion/data-index-ingestion-kafka-processor/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/processor/persistence/WorkflowPersistence.java @@ -17,6 +17,7 @@ import java.sql.SQLException; import java.sql.Types; import java.time.ZonedDateTime; +import java.util.List; import java.util.Optional; @Unremovable @@ -58,6 +59,31 @@ public void persist(WorkflowInstance event) throws SQLException { } } + public void persistBatch(List events) throws SQLException { + if (events == null || events.isEmpty()) { + return; + } + + try (Connection conn = dataSource.getConnection(); + PreparedStatement stmt = conn.prepareStatement(insertWorkflowUpsert)) { + + conn.setAutoCommit(false); + + try { + for (WorkflowInstance event : events) { + setWorkflowParameters(stmt, event); + stmt.addBatch(); + } + + stmt.executeBatch(); + conn.commit(); + } catch (SQLException e) { + conn.rollback(); + throw e; + } + } + } + private void setWorkflowParameters(PreparedStatement stmt, WorkflowInstance event) throws SQLException { stmt.setString(1, event.getId()); stmt.setString(2, event.getNamespace()); @@ -105,4 +131,4 @@ private String toJsonString(JsonNode node) { } } -} \ No newline at end of file +} diff --git a/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/service/KafkaLifecycleConsumer.java b/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/service/KafkaLifecycleConsumer.java index 5ad0d66e5..0ce368cf4 100644 --- a/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/service/KafkaLifecycleConsumer.java +++ b/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/java/org/kubesmarts/logic/dataindex/ingestion/kafka/service/KafkaLifecycleConsumer.java @@ -21,60 +21,164 @@ import io.cloudevents.jackson.JsonCloudEventData; import io.serverlessworkflow.impl.lifecycle.ce.TaskCEData; import io.serverlessworkflow.impl.lifecycle.ce.WorkflowCEData; +import io.smallrye.reactive.messaging.MutinyEmitter; +import io.smallrye.reactive.messaging.kafka.api.OutgoingKafkaRecordMetadata; import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.common.header.internals.RecordHeaders; +import org.eclipse.microprofile.config.inject.ConfigProperty; +import org.eclipse.microprofile.reactive.messaging.Channel; import org.eclipse.microprofile.reactive.messaging.Incoming; +import org.eclipse.microprofile.reactive.messaging.Message; +import org.eclipse.microprofile.reactive.messaging.Metadata; import org.kubesmarts.logic.dataindex.ingestion.kafka.processor.EventProcessor; import org.kubesmarts.logic.dataindex.ingestion.kafka.processor.ProcessEventFailedException; +import org.kubesmarts.logic.dataindex.ingestion.kafka.processor.WorkflowEventProcessor; +import org.kubesmarts.logic.dataindex.ingestion.kafka.processor.TaskExecutionProcessor; + import org.kubesmarts.logic.dataindex.model.LifecycleEventUtils; import org.kubesmarts.logic.dataindex.model.TaskExecution; import org.kubesmarts.logic.dataindex.model.WorkflowInstance; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; + + @ApplicationScoped public class KafkaLifecycleConsumer { private static final Logger log = LoggerFactory.getLogger(KafkaLifecycleConsumer.class); + private final WorkflowEventProcessor workflowEventProcessor; + private final TaskExecutionProcessor taskExecutionProcessor; + + @Inject + public KafkaLifecycleConsumer(WorkflowEventProcessor workflowEventProcessor, + TaskExecutionProcessor taskExecutionProcessor) + { + this.workflowEventProcessor = workflowEventProcessor; + this.taskExecutionProcessor = taskExecutionProcessor; + } + @Inject ObjectMapper jackson; @Inject - EventProcessor workflowEventProcessor; + @Channel("data-index-events-dlq") + MutinyEmitter deadLetterEmitter; - @Inject - EventProcessor taskExecutionProcessor; + @ConfigProperty(name = "data-index.ingestion.db-batch-size", defaultValue = "1000") + int dbBatchSize; @Incoming("data-index-events") - public void consumeLifecycleEvent(ConsumerRecord record) { - + public CompletionStage consumeLifecycleEvent(Message> records) { + List> deadLetterSends = new ArrayList<>(); + + List workflowInstances = new ArrayList<>(); + List taskExecutions = new ArrayList<>(); + + for (ConsumerRecord record : records.getPayload()) { + try { + CloudEvent cloudEvent = validateCloudEvent(record); + + JsonCloudEventData cloudEventData = (JsonCloudEventData) cloudEvent.getData(); + if (cloudEventData == null || cloudEventData.getNode() == null) { + throw new IllegalArgumentException("The CloudEvent data node consumed at offset %s from partition %s is null or empty." + .formatted(record.offset(), record.partition())); + } + + Class eventClass = LifecycleEventUtils.getEventClass(cloudEvent.getType()); + Object data = jackson.convertValue(cloudEventData.getNode(), eventClass); + + if (data instanceof TaskCEData taskData) { + taskExecutions.add(mapTaskEvent(cloudEvent, taskData)); + } else if (data instanceof WorkflowCEData workflowData) { + workflowInstances.add(mapWorkflowEvent(cloudEvent, workflowData)); + } + } catch (Exception e) { + log.error("Failed to consume the record from Kafka at offset '{}' from partition '{}'. Routing to dead-letter queue.", + record.offset(), record.partition(), e); + deadLetterSends.add(sendToDeadLetterQueue(record, e)); + } + } + try { - CloudEvent cloudEvent = validateCloudEvent(record); - - JsonCloudEventData cloudEventData = (JsonCloudEventData) cloudEvent.getData(); - if (cloudEventData == null || cloudEventData.getNode() == null) { - throw new IllegalArgumentException("The CloudEvent data node consumed at offset %s from partition %s is null or empty." - .formatted(record.offset(), record.partition())); + log.info("Mapped Kafka batch: workflows={}, tasks={}", + workflowInstances.size(), taskExecutions.size()); + + for (List chunk : partition(workflowInstances, dbBatchSize)) { + workflowEventProcessor.processBatch(chunk); } - - Class eventClass = LifecycleEventUtils.getEventClass(cloudEvent.getType()); - Object data = jackson.convertValue(cloudEventData.getNode(), eventClass); - - if (data instanceof TaskCEData taskData) { - handleTaskEvent(cloudEvent, taskData); - } else if (data instanceof WorkflowCEData workflowData) { - handleWorkflowEvent(cloudEvent, workflowData); - } else { - throw new IllegalArgumentException("Unsupported event type '%s' consumed at offset %s from partition %s." - .formatted(cloudEvent.getType(), record.offset(), record.partition())); + + for (List chunk : partition(taskExecutions, dbBatchSize)) { + taskExecutionProcessor.processBatch(chunk); } } catch (Exception e) { - log.error("Failed to consume the record from Kafka at offset '{}' from partition '{}'.", record.offset(), record.partition(), e); - throw new ProcessEventFailedException("Failed to consume Kafka record at offset %s from partition %s".formatted( - record.offset(), record.partition()), e); + log.error("Failed to persist Kafka batch. Nacking the batch so it can be retried.", e); + return records.nack(e); + } + + CompletableFuture[] pending = deadLetterSends.stream() + .map(CompletionStage::toCompletableFuture) + .toArray(CompletableFuture[]::new); + + return CompletableFuture.allOf(pending).thenCompose(ignored -> records.ack()); + } + + private static List> partition(List list, int size) { + List> chunks = new ArrayList<>(); + + for (int i = 0; i < list.size(); i += size) { + chunks.add(list.subList(i, Math.min(i + size, list.size()))); } + + return chunks; + } + + private TaskExecution mapTaskEvent(CloudEvent cloudEvent, TaskCEData data) { + try { + return Mapper.mapTaskExecutionEvent(cloudEvent, data, jackson); + } catch (Exception e) { + log.error("Error while mapping CloudEvent (task) with ID: {}", cloudEvent.getId(), e); + throw new ProcessEventFailedException("Failed to map CloudEvent with ID: " + cloudEvent.getId(), e); + } + } + + private WorkflowInstance mapWorkflowEvent(CloudEvent cloudEvent, WorkflowCEData data) { + try { + return Mapper.mapWorkflowInstanceEvent(cloudEvent, data, jackson); + } catch (Exception e) { + log.error("Error while mapping CloudEvent (workflow) with ID: {}", data.getName(), e); + throw new ProcessEventFailedException("Failed to map CloudEvent with ID: " + data.getName(), e); + } + } + + private CompletionStage sendToDeadLetterQueue(ConsumerRecord record, Exception cause) { + RecordHeaders headers = new RecordHeaders(); + headers.add("dead-letter-reason", bytes(cause.getMessage() != null ? cause.getMessage() : cause.toString())); + headers.add("dead-letter-cause", bytes(cause.getClass().getName())); + headers.add("dead-letter-original-topic", bytes(record.topic())); + headers.add("dead-letter-original-partition", bytes(Integer.toString(record.partition()))); + headers.add("dead-letter-original-offset", bytes(Long.toString(record.offset()))); + + OutgoingKafkaRecordMetadata metadata = OutgoingKafkaRecordMetadata.builder() + .withKey(record.key()) + .withHeaders(headers) + .build(); + + return deadLetterEmitter.sendMessage(Message.of(record.value(), Metadata.of(metadata))) + .subscribeAsCompletionStage(); + } + + private static byte[] bytes(String value) { + return value.getBytes(StandardCharsets.UTF_8); } private CloudEvent validateCloudEvent(ConsumerRecord record) throws JsonProcessingException { diff --git a/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/resources/application.properties b/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/resources/application.properties index 87e21b118..4ed6747fa 100644 --- a/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/resources/application.properties +++ b/data-index/data-index-ingestion/data-index-ingestion-kafka-service/src/main/resources/application.properties @@ -19,6 +19,7 @@ quarkus.container-image.tag=999-SNAPSHOT # mp messaging mp.messaging.incoming.data-index-events.connector=smallrye-kafka mp.messaging.incoming.data-index-events.topic=flow-lifecycle-out +mp.messaging.incoming.data-index-events.batch=true mp.messaging.incoming.data-index-events.group.id=data-index-ingestion mp.messaging.incoming.data-index-events.auto.offset.reset=earliest mp.messaging.incoming.data-index-events.retry-attempts=2 @@ -27,9 +28,14 @@ mp.messaging.incoming.data-index-events.key.deserializer=org.apache.kafka.common mp.messaging.incoming.data-index-events.health-enabled=true mp.messaging.incoming.data-index-events.health-readiness-enabled=true -mp.messaging.incoming.data-index-events.failure-strategy=dead-letter-queue -mp.messaging.incoming.data-index-events.dead-letter-queue.topic=data-index-events-dlq -mp.messaging.incoming.data-index-events.dead-letter-queue.key.serializer=org.apache.kafka.common.serialization.StringSerializer -mp.messaging.incoming.data-index-events.dead-letter-queue.value.serializer=org.apache.kafka.common.serialization.StringSerializer +# Per-record failures are dead-lettered manually by KafkaLifecycleConsumer (batch ack/nack are +# whole-batch operations). 'fail' only triggers if a dead-letter write itself fails, so the batch +# is not committed and the failed record is not silently dropped. +mp.messaging.incoming.data-index-events.failure-strategy=fail + +mp.messaging.outgoing.data-index-events-dlq.connector=smallrye-kafka +mp.messaging.outgoing.data-index-events-dlq.topic=data-index-events-dlq +mp.messaging.outgoing.data-index-events-dlq.key.serializer=org.apache.kafka.common.serialization.StringSerializer +mp.messaging.outgoing.data-index-events-dlq.value.serializer=org.apache.kafka.common.serialization.StringSerializer quarkus.native.resources.includes=task-instance-upsert.sql,task-placeholder-workflow-insert.sql,workflow-instance-upsert.sql \ No newline at end of file