From d057900c3dc70a15e7ee0673bf0426877ffbb98e Mon Sep 17 00:00:00 2001 From: Karl F Walkow Date: Fri, 17 Jul 2026 13:07:25 +0200 Subject: [PATCH 1/3] feat: Make throughput logging configurable --- src/main/scala/de/otto/anthology/App.scala | 3 ++- src/main/scala/de/otto/anthology/AppWorkflow.scala | 5 +++-- src/main/scala/de/otto/anthology/DomainSources.scala | 8 +++++--- src/main/scala/de/otto/anthology/KafkaSink.scala | 4 ++-- .../scala/de/otto/anthology/config/DomainConfig.scala | 3 ++- 5 files changed, 14 insertions(+), 9 deletions(-) diff --git a/src/main/scala/de/otto/anthology/App.scala b/src/main/scala/de/otto/anthology/App.scala index 19a8a62..5bb0d19 100644 --- a/src/main/scala/de/otto/anthology/App.scala +++ b/src/main/scala/de/otto/anthology/App.scala @@ -110,7 +110,8 @@ object App extends OxApp, LazyLogging: store, clusterSettings, kafkaConsumers, - config.parallelism + config.parallelism, + config.domain.logThroughput ) par(startHttpServer(), startAppWorkflow()).discard ExitCode.Success diff --git a/src/main/scala/de/otto/anthology/AppWorkflow.scala b/src/main/scala/de/otto/anthology/AppWorkflow.scala index 30a84c0..6708ac1 100644 --- a/src/main/scala/de/otto/anthology/AppWorkflow.scala +++ b/src/main/scala/de/otto/anthology/AppWorkflow.scala @@ -44,9 +44,10 @@ object AppWorkflow: stateStore: StateStore, clusterSettings: Map[ClusterName, KafkaClusterSettings], kafkaConsumers: ConsumerMap, - parallelism: Parallelism + parallelism: Parallelism, + logDomainThroughput: Option[Boolean] )(using Ox): Unit = - DomainSources(channelConfigs, kafkaConsumers) + DomainSources(channelConfigs, kafkaConsumers, logDomainThroughput) .buffer() .filterDomainMessages(channelConfigs, parallelism) .buffer() diff --git a/src/main/scala/de/otto/anthology/DomainSources.scala b/src/main/scala/de/otto/anthology/DomainSources.scala index 01d0650..080aff8 100644 --- a/src/main/scala/de/otto/anthology/DomainSources.scala +++ b/src/main/scala/de/otto/anthology/DomainSources.scala @@ -14,14 +14,16 @@ object DomainSources: def apply( configs: ChannelConfigs, - consumers: ConsumerMap + consumers: ConsumerMap, + logThroughput: Option[Boolean] )(using Ox): Flow[(Option[(QualifiedMessageId, Option[Message])], Passthrough)] = configs.channels .map: config => // for now we go with consumer name == domain name val consumerName = ConsumerName(config.name.toString) - DomainSource( + val src = DomainSource( config, KafkaSourceSettings(config.kafka, consumerName, consumers(consumerName)) - ).count(s"domain source ${config.name.toString}") + ) + if logThroughput.getOrElse(false) then src.count(s"domain source ${config.name.toString}") else src .mergeFair() diff --git a/src/main/scala/de/otto/anthology/KafkaSink.scala b/src/main/scala/de/otto/anthology/KafkaSink.scala index 4a20157..212b069 100644 --- a/src/main/scala/de/otto/anthology/KafkaSink.scala +++ b/src/main/scala/de/otto/anthology/KafkaSink.scala @@ -97,8 +97,8 @@ object KafkaSink extends LazyLogging: producerRecords.foreach(publishChannel.send) offsets.foreach(committerChannel.send) .mapStateful(0): (cnt, _) => - if cnt % 100 == 0 then - logger.info("Published and committed 100 batches of codomain messages to Kafka") + if cnt % 1000 == 0 then + logger.info("Published and committed 1000 batches of codomain messages to Kafka") (cnt + 1, ()) .runDrain() diff --git a/src/main/scala/de/otto/anthology/config/DomainConfig.scala b/src/main/scala/de/otto/anthology/config/DomainConfig.scala index 09206e5..dfe60a3 100644 --- a/src/main/scala/de/otto/anthology/config/DomainConfig.scala +++ b/src/main/scala/de/otto/anthology/config/DomainConfig.scala @@ -3,5 +3,6 @@ package de.otto.anthology.config import de.otto.anthology.ChannelName import pureconfig.ConfigReader -case class DomainConfig(channels: Seq[ChannelConfig], relations: Seq[RelationConfig]) derives ConfigReader: +case class DomainConfig(channels: Seq[ChannelConfig], relations: Seq[RelationConfig], logThroughput: Option[Boolean]) + derives ConfigReader: val channelsByName: Map[ChannelName, ChannelConfig] = channels.map(c => (c.name, c)).toMap From 04ed81fe1a5966a6b4cca4d37b6ab729e3e93fb4 Mon Sep 17 00:00:00 2001 From: Karl F Walkow Date: Fri, 17 Jul 2026 13:11:17 +0200 Subject: [PATCH 2/3] feat: Reduce stream buffers --- src/main/scala/de/otto/anthology/AppWorkflow.scala | 5 ----- 1 file changed, 5 deletions(-) diff --git a/src/main/scala/de/otto/anthology/AppWorkflow.scala b/src/main/scala/de/otto/anthology/AppWorkflow.scala index 6708ac1..5bcacd6 100644 --- a/src/main/scala/de/otto/anthology/AppWorkflow.scala +++ b/src/main/scala/de/otto/anthology/AppWorkflow.scala @@ -48,11 +48,8 @@ object AppWorkflow: logDomainThroughput: Option[Boolean] )(using Ox): Unit = DomainSources(channelConfigs, kafkaConsumers, logDomainThroughput) - .buffer() .filterDomainMessages(channelConfigs, parallelism) - .buffer() .transformDomainMessageIds(channelConfigs) - .buffer() .transformDomainMessages(channelConfigs, parallelism) .buffer() .persistDomainMessages(stateStore) @@ -66,13 +63,11 @@ object AppWorkflow: .inlineDomainMessages(relationConfigs, stateStore, parallelism) .buffer() .filterCodomainMessages(codomainConfig.filtering, parallelism) - .buffer() .transformCodomainMessages(codomainConfig.transformation, parallelism) .buffer() .persistCodomainMessages(stateStore, parallelism) .buffer() .propagateHeaders(codomainConfig.headerPropagationConfigs) - .buffer() .emitCodomainMessages( KafkaSinkSettings( codomainConfig.kafka, From 0c2578570f874775d38d496fbd8f14d700b1625d Mon Sep 17 00:00:00 2001 From: Karl F Walkow Date: Fri, 17 Jul 2026 13:23:48 +0200 Subject: [PATCH 3/3] feat: skip batching when batch size is null --- .../CodomainDeduplicationStage.scala | 36 ++++++++++--------- 1 file changed, 20 insertions(+), 16 deletions(-) diff --git a/src/main/scala/de/otto/anthology/CodomainDeduplicationStage.scala b/src/main/scala/de/otto/anthology/CodomainDeduplicationStage.scala index 7f6ee55..4e524fa 100644 --- a/src/main/scala/de/otto/anthology/CodomainDeduplicationStage.scala +++ b/src/main/scala/de/otto/anthology/CodomainDeduplicationStage.scala @@ -19,21 +19,25 @@ object CodomainDeduplicationStage: configOpt: Option[CodomainDeduplicationConfig] ): Flow[(Seq[(QualifiedMessageId, Seq[MessageId])], Seq[Passthrough])] = val config = configOpt.getOrElse(defaultConfig) - in - .groupedWithin(config.batchSize, config.batchingDuration) - .map: batches => - val passthroughs: ListBuffer[Passthrough] = ListBuffer.empty - val deduplicationMap: MutableMap[QualifiedMessageId, Set[MessageId]] = MutableMap.empty - batches.foreach: batch => - batch._1 match - case None => - () - case Some(domainMessageId, codomainMessageIds) => - deduplicationMap.updateWith(domainMessageId): - case None => Some(Set.empty ++ codomainMessageIds) - case Some(curCodomainMessageIds) => - Some(curCodomainMessageIds ++ codomainMessageIds) - passthroughs += batch._2 - (deduplicationMap.map((k, v) => (k, v.toSeq)).toSeq, passthroughs.sorted.toSeq) + if config.batchSize == 1 then + in.map: (payload, passthrough) => + (payload.toSeq.map(e => (e._1, e._2.toSeq)), Seq(passthrough)) + else + in + .groupedWithin(config.batchSize, config.batchingDuration) + .map: batches => + val passthroughs: ListBuffer[Passthrough] = ListBuffer.empty + val deduplicationMap: MutableMap[QualifiedMessageId, Set[MessageId]] = MutableMap.empty + batches.foreach: batch => + batch._1 match + case None => + () + case Some(domainMessageId, codomainMessageIds) => + deduplicationMap.updateWith(domainMessageId): + case None => Some(Set.empty ++ codomainMessageIds) + case Some(curCodomainMessageIds) => + Some(curCodomainMessageIds ++ codomainMessageIds) + passthroughs += batch._2 + (deduplicationMap.map((k, v) => (k, v.toSeq)).toSeq, passthroughs.sorted.toSeq) case class CodomainDeduplicationConfig(batchSize: Int, batchingDuration: FiniteDuration) derives ConfigReader