Saturday, 19 September 2026

[3 of 3] I benchmarked five Scala libraries. Which code would I rather read?

I started with synthetic benchmarks in part 1. Adding socket I/O in part 2 sent me into Kyo's scheduler. Along the way, Ox and Gears caught my attention. Their TCP throughput was close to cats-effect, they allocated much less, and I found their direct-style code easier to read.

For this last post, I added plain Loom and checked whether lower allocation also meant less memory. The results still don't give me one choice that wins everywhere.

Five runtimes, different answers

These are selected results from the new run, in operations per second. Higher is better. The ± values are JMH's 99.9% confidence-interval half-widths; compare libraries within each column.

Throughput, including runtime entry and task management
Runtime8 workersSpawn/joinBlocking TCPCallback TCP
cats-effect3,060 ± 921,894 ± 85145.10 ± 9.03141.53 ± 3.07
Kyo1,399 ± 265,355 ± 1044.95 ± 0.03139.99 ± 1.73
Loom / JDK3,420 ± 71558 ± 43143.49 ± 5.65137.69 ± 3.44
Ox3,580 ± 230494 ± 51140.29 ± 6.62136.78 ± 3.48
Gears3,214 ± 50377 ± 26140.28 ± 9.81139.25 ± 2.57

The eight-worker case processes 4,096 values, with 64 arithmetic rounds per value. Spawn/join creates a child, waits for it to finish, and repeats 1,000 times. Each TCP operation completes 256 exchanges over 64 persistent loopback connections. A separate server requests a 1 ms delay per response.

Kyo leads the tiny-task spawn/join case, while the callback TCP intervals overlap across all five. The slow Kyo blocking result uses defaults. In part 2, flushing and worker tuning brought it close to CE. That tuned result belongs to the earlier experiment.

All five ran every case shown here. The Loom row uses JDK virtual threads and primitives. The suspended effect-chain tests from part 1 have no equivalent Loom, Ox or Gears benchmark: replacing the chain with a direct loop changes what we measure. I haven't compared cancellation or resource-safety guarantees.

Less allocation didn't answer the memory question

I kept 100,000 tasks waiting, each holding a 1 KiB payload, and measured live heap after full GC and macOS process footprint. I also sampled memory during blocking TCP work at 64 connections. The numbers are medians from three fresh JVMs per case, in MiB. TCP measurements cover only the client.

Memory consumption: lower is better
RuntimeWaiting tasks: live heapWaiting tasks: process footprintActive TCP: process footprint
cats-effect191.0418.5152.5
KyoNot measured
Loom / JDK228.7390.090.3
OxNot measured
GearsNot measured

CE used less live heap for these waiting tasks. Loom had the smaller process footprint, which includes memory beyond live heap objects. So far, the memory harness covers only CE and Loom. I haven't measured Kyo, Ox or Gears here; the blanks don't indicate missing library support.

Both studies ran on 19 September 2026 on one M3 Max with JDK 25.0.3 and Scala 3.8.4: CE 3.7.1, Kyo 1.0.0-RC6, Ox 1.0.6 and Gears 0.3.1. Throughput used three forks, three 1-second warmups and three 1-second measurements, with a fixed 2 GiB G1 heap. Memory used G1 with a 64 MiB initial and 2 GiB maximum heap. These are short runs of these particular implementations on one machine.

What I want to maintain

I still enjoy writing with cats-effect and Kyo. After moving into Tagless Final, though, I sometimes find myself reading through a lot of ceremony to get to the business logic. Making every service generic in F[_] is a choice I can reconsider while still using CE.

That makes me more interested in direct-style concurrency. I often find its control flow easier to follow, even though writing effects is more fun for me. The benchmark results give me reasons to consider both.

AI makes me reconsider what I value

AI tools handled most of my earlier effect-system migration. If I spend more time reviewing generated code and less time writing it, how much should my enjoyment of writing a particular style count?

That even puts Java back on the table for me. Boilerplate bothers me less when I'm not typing all of it, but I still have to read and maintain it. I could live with more verbose code if I found it easier to follow what runs, how it fails and who owns a resource.

Effects and types can help me review code too. A version with fewer combinators might still hide a cancellation or cleanup bug.

I don't have one library to recommend for every project. I'd choose around the workload, the guarantees I need and the integrations available. Among the options that fit, I'd give more weight to what my team can comfortably read and maintain, whoever wrote it.

Sunday, 13 September 2026

[2 of 3] Blocking I/O made Kyo 30x slower. Then I tried flush().

I spent a week giving Kyo a chance in a real service. AI coding tools handled most of the effect-system migration.

After the first synthetic comparison, I also wanted to see what happened when the programs had to wait for actual socket I/O. The numbers below come from a separate TCP benchmark.

I added Gears and Ox too. Same TCP exchanges, connection limits and ordered results. The nonblocking numbers were close. Blocking exposed a much bigger difference.

Kyo led cats-effect by about 12% for nonblocking I/O at concurrency 8. At 64, it trailed by about 3%. That is already too mixed for a blanket claim that Kyo is faster. The earlier synthetic wins still describe those particular operations.

For blocking I/O at concurrency 64, cats-effect completed about 149 batches per second. Default Kyo managed about 5.

Each batch contained 256 requests over 64 persistent connections. A separate server JVM requested a 1 ms sleep before each response. These measurements used cats-effect 3.7.1 and Kyo 1.0.0-RC6 on one Mac with JDK 25.0.3, recorded on 6 September 2026.

Blocking configurationBatches/s, higher is better
cats-effect default149.00 ± 2.93
Kyo default4.96 ± 0.02
Kyo with flush before every exchange81.86 ± 25.62
Kyo with flush and scheduler tuning144.89 ± 3.60

Three JVM forks per case. The ± values are JMH's 99.9% confidence-interval half-widths. Flush alone varied substantially; the tuned Kyo and default CE intervals overlap.

The source investigation pointed to queued children collecting on a few workers, slow adaptation to short blocking calls, and a placement scan that could miss idle workers.

Calling the scheduler's flush hook before each exchange redistributed queued work:

Sync.defer {
    kyo.scheduler.Scheduler.get.flush()
    io.blocking(lane, index)
}

That is the blocking branch from the benchmark. Flush made a huge difference, but reaching CE's range also required setting coreWorkers, minWorkers, maxWorkers and scheduleStride to 64. Those are experimental settings for this workload. No library code changed.

Flush does not move the blocking call to another thread or create more workers. It gives queued tasks another chance to run elsewhere. We called it before every exchange; we have not tested whether once per worker would be enough.

The manual step bothers me. Kyo supports blocking through Sync.defer, yet getting good performance here required knowing about a scheduler hook. CE handled the same workload well with IO.blocking and its default runtime.

Kyo's Finagle integration already calls flush inside its blocking hook. When that integration is enabled and the call goes through the hook, it performs the flush automatically. That does not cover every arbitrary blocking call, and it does not apply our worker tuning. We did not benchmark Finagle.

Putting flush in every Sync.defer would also affect cheap side effects. My reading is that this could add scheduling overhead; I have no maintainer confirmation of that design rationale. A dedicated blocking helper seems worth exploring.

Gears and Ox landed around 139 blocking batches/s at concurrency 64 and allocated much less than CE. Their direct style is interesting even without a throughput win: ordinary loops, conditions and local helpers around I/O, with fewer effect combinators. That may matter more in everyday code than a small benchmark lead.

Kyo still has my attention. I also want to know how much scheduler knowledge an application will need.

Full tables, allocation, CPU and raw measurements, plus methods and reproduction. The suite passed 438 correctness checks. This is one loopback workload; it does not establish production performance or equivalent cancellation guarantees.

Saturday, 5 September 2026

[1 of 3 ]Kyo vs cats-effect: promising numbers and one big surprise

Kyo vs cats-effect: promising numbers and one big surprise

I find Kyo one of the most interesting projects in Scala right now. I get the feeling that some of the complexity we accept might not be necessary.

That feeling is easy to get excited about. I wanted to see some numbers too. This is the first post in a series about Kyo and cats-effect, starting with a small synthetic war.

I compared cats-effect 3.7.1 with Kyo 1.0.0-RC6, using fs2 3.13.0 on the cats-effect side for streaming.

Getting a fair comparison took more work than I expected. Two methods can look equivalent while doing different things with batches, buffers, or failures.

The suite passed 984 correctness checks. I ran JMH with three JVM forks and identical heap settings on one Mac with JDK 25.

Both sides use the same worker strategy and stream batch boundaries. The success-collection row uses parallel attempts and filtering, rather than native Async.gather. The stream rows do not directly compare parEvalMap with mapPar.

These are selected results with very little work per element. Ratios compare mean throughput; CE means cats-effect.

ConstructionHigher throughput
Uncontended permit acquisition/releaseKyo ~5.02×
Queue-backed stream batchesKyo ~3.70×
Sequential fiber spawn/joinKyo ~2.75×
Parallel stream batchesKyo ~1.93×
Integer queue, one producer and consumerKyo ~1.45×
Bounded workersKyo ~1.13×
Parallel attempt and success collectionConfidence intervals overlap
Left-associated suspended binds, depth 10,000CE ~1,895×

The permit, queue, and fiber results make Kyo worth a closer look. Around 5× the throughput for uncontended permits and 3.7× for queue-backed stream batches are encouraging numbers, even for these small workloads.

A little CPU work changed the results too. Bounded workers went from a small Kyo lead to CE ahead by about 2.22×. Kyo’s lead in parallel stream batches shrank to about 1.11×.

Then there’s the last row. CE handled the 10,000-step left-associated bind chain at about 1,895× Kyo’s throughput. Kyo allocated about 1.6 GB per complete chain, versus roughly 943 KB for CE. Making the chain ten times longer increased Kyo’s allocation about a hundredfold. That points to quadratic scaling in this case, and it needs a closer look.

These results make me more interested in trying Kyo. There’s enough here to justify a small application experiment, and a specific weakness to investigate along the way.

This is one machine and a set of small workloads. Cancellation, resource safety, and real I/O still need their own testing.

Source, full results, confidence intervals, and reproduction commands.

If you spot an unfair comparison, please point me to the code. I’d like to get it right.

Thursday, 28 February 2019

Aux pattern in 5 minutes

Every-time using everytime shapeless.Generic that returns shapeless.Generic.Aux
I was wondering why the construction
type Aux[T, Repr0] = Generic[T] { type Repr = Repr0 }
is needed?

If you are impassioned and want to save a time: AUX pattern is wrapping abstract type member into generic parametrisation to solve compiler limitations.

While the Scala's abstract type members are preferred to generic parametrization sometimes there are the problems with making the code compilable.
Lets revisit Type level programming for Streams and start from the case when Aux pattern isn’t needed:

import cats.Show

implicit val requestFunc = new Function[Int, String] {
  override def apply(in: Int): String = s"hello $in"
}
def process[Req, Resp](request: Req)(implicit func: Function[Req, Resp], s: Show[Resp]): String 
  = func(request).show

Our process function just applies the func to request and then using Show type class to represent the result value. Show and func are implicitly injected according to Req and Resp types.

Lets rewrite using abstract type member:

trait Request[Req] {
  type Resp
  def response: Resp
}

implicit def string2Sting = new Request[Int] {
  type Resp = String
  override def response: String = "Hello"
}

def process2[T](value: T)(implicit request: Request[T], m: Show[request.Resp]): String = 
 request.response.show

Method process2 isn’t compilable because we are trying parametrise Show with request.Response type:
error: illegal dependent method type: parameter may only be referenced in a subsequent parameter section

To fix this we introduce proxy type (Aux pattern).

type Aux[T2, B2] = Request[T2] {type B = B2}

And changing the method to:

def process3[Req, Resp](request: Req)(implicit aux: Aux[Req, Resp], m: Show[Resp]): String = 
aux.response.show

As a result we create Aux type that wraps abstract type member into generic to overcome Scala's compiler limitation!

Tuesday, 19 February 2019

Switch career to machine learning specialist

There is continues hype around AI. As this hype stream is still far from been over and predictively is going to grow in time it's natural for the software engineers to pry what is happening in the area of AI Engineer Jobs.

A lot of us been career switchers - got non IT education math or engineering and end up in IT industry as software engineers. Exceptional professions like 3D engine developers etc always have been demanding to heavy math background.

What is about ML - there are many courses/books like learn it in N days. Career switchers are used to get an experience driven into new background with consistent feeling lick of knowledge.

Does this skill help with diving into ML? My opinion it isn't - before starting ML courses it's mandatory to learn/refresh match background - here is graph for studying ML from scratch:




Tuesday, 11 December 2018

Hylomorphism in 1 minute

A Hylomorphism is  Catamorphism compose to Anamorphism function.

It's presented as 

Hylomorphism = catamorphism(anamorphism(x))
 or
Hylomorphism = fold(unfold(x))

Basically it constructs (unfold) complex type (like trees, lists) and destructs (folds) back into representing value.

For example to get factorial from N we can 
a) unfold N to list of (n),(n-1),..,0
b) fold with prod function the list from previous step.

And possible Scala example implementing function for getting factorial from N with help of previously introduced list's Catamorphism and  Anamorphism is:

Anamorphism in 1 minute


An Anamorphism (from the Greek ἀνά "upwards" and μορφή "form, shape) over a type T is a function of type U → T that destructs the U in some way into a number of components. On the components that are of type U, it applies recursion to convert theme to T 's. After that it combines the recursive results with the other components to a T , by means of the constructor functions of T .

For example if we want to build List(N, N-1, …, 1) from N we would use anamorphism from Int to List[Int]. It’s the opposite operation to fold - unfold.


Catamorphism in 2 minutes


A Catamorphism (κατά "downwards" and μορφή "form, shape") on a type T is a function of type T → U that destruct and object of type T according to the structure of T, calls itself recursively of any components of T that are also of type T and combines this recursive result with the remaining components of T to a U.

And it's just an official name of fold/reduce on higher kinded types. 

For example getting the prod from list of integers is Catamorphism on the list:


Catamorphism in programming can be used to fold complex structures (lists, tress, sets)  into their representation via different type. As an example list catamorphism can be described as


And if we want to fold the list into prod from elements we would use:


Tuesday, 4 December 2018

Isomorphism in 3 minutes



An Isomorphism (from the Ancient Greek: ἴσος isos "equal", and μορφή morphe "form" or "shape") is a homomorphism or morphism (i.e. a mathematical mapping) that can be reversed by an inverse morphism. Two mathematical objects are isomorphic if an isomorphism exists between them.
Isomorphism between the types means that we can convert T  →  U and then U → T lossless.

For example Int to String conversion is sort of isomorphism, but not all the possible values of String can be converted to Int. For example when we try to convert “Assa” into Int we get an Exception. If we want to use identity element for example 0 for any non-number String - Int2String won’t be Isomorphic anymore.


The real isomorphic mapping from String to Int can be done via co-algebra (list’s co-algebra):


You’ve probably generated the tons of Unit tests for  JSON → String and String → JSON isomorphism prove.

Singleton types are always isomorphic to itself, for example type Unit has only one set’s member Unit or () and it’s always isomorphic to itself. Scala compiler has a special flag -Ywarn-value-discard that checks method returns type may break isomorphism rule (make sense only for non side-effect calls).



Currying functions are isomorphic to each other:

Function1 in Scala doesn’t have the curry method and that is why curry isn’t isomorphic for all the Scala’s functions - but it could be if curry from Function1 returned Function1[T, Function1[T1, U]].


Isomorphism as applied to web development means been able to render pages on the both server and client sides. It’s used in context of NodeJS ecosystem because of been able to reuse libraries/frameworks in backend and frontend.



Friday, 2 November 2018

JavaScript: the legacy and pain in tc39

"You Might Not Need TypeScript." © Is a bad statement.
JavaScript is evolving and expanding with the features but.

Let's look into the latest language feature been implemented in Chrome browser: Class Fields.
There is a link to Technical Committee 39 developing the specification itself. According to github discussion has been started 4+ years ago, we should expect some mature and important language feature.

The ECMAScript proposal “Class Fields” is about fields declarations for classes, for example to declare public instance property we can use:
Easy isn't it? I would expect it's a sugar for:
Actually it's not, the code behaves unpredictable when is used in class inheritance. If we define Parent class with some property and extend it in Child class with overriding the property - it behaves as expected:
But let's refactor parent class and keep expected backward compatibility:
Oops, it doesn't work as expected. It's because of Object.defineProperty is used over normal set operation. It becoming the major difference because property addition through set creates properties which show up during property enumeration, whose values may be changed, and which may be deleted. Object.defineProperty is defining them lazily, that brings to the cases when it can be shadowed by set operations. Here are more details: [[Set]] vs [[DefineOwnProperty]].

Update: Babel 7.1.0 changes default behaviour from [[Set]].

The legacy doesn't allow language to grow in fair and proper way, I have a feeling that JavaScript is moving in the same direction as a Perl done, but JavaScript is protected from been abandonded one's days because it's a web language.


Tuesday, 30 October 2018

#2018 Trends



It's almost November but feels like a winter, I've tried to create cheat sheet with been trendy buzzwords/streams this year I had to face or work with. Did it for myself to read/refresh books/blogs during the upcoming vacation but might somebody will find it useful to be shared:
  • Java from Oracle isn't free anymore
  • Java as a language is lagging behind all others, Lombok and Spring Roo aren't demanded solutions. I still bet to Scala
  • Kotlin got popularity, finally 1.3 has been released
  • Rust 1.3 has been released
  • DevSecOps over DevOps, it's mandatory to invest into security
  • Prometheus over Graphite
  • Microservices is still the buzzword
  • Docker vs Podman, project CRI-O, OCID (Open Container Initiative Daemon) - well Ops in nova days is something that is quickly changing
  • IBM got Redhat, are there any other open-source independent whales?
  • Clouds becoming private for intermediate size companied
  • Ansible over Chef and Puppet
  • Python is great again
  • AI buzzword over data analysis (k-nearest neighbors classifier is accepted as magic spell)
  • Data Scientist and Data Engineer are demanded, to get outside of Hadoop, Spark, Kudu, Impala, Kafka, Ignite, Geode etc is a separate profession
  • Engineering Culture over Agile to address Conway's law
  • Vue, React are gaining popularity
  •  Kafka as a glue for everything. Got the competitor nats.io streams but for streaming it's still astable (by the way implemented in Go)
  • Yaml is for configuration
  • GraphQL is preferred to REST
  • Cassandra is preferred to MongoDB
  • GraalVM: is too a day before the fair to start looking into
  • JVM: still there is a demand for green threads

Tuesday, 9 October 2018

Scala Сheat sheet: Context Bound of multi typed kind

Because of rarely using this functionality every time has to find out how to use multi type for context bound.

For example there is class extending Function1:

abstract class Mapping[In, Out] extends (In => Out)

If we want to use Mapping[In, String] as a Context Bound for other class we should use type Lambda (type-level Lambda):

abstract class StringRequestResponse[In: ({type M[x] = Mapping[x, String]})#M]

In case of return should be parametrised as well:

abstract class RequestResponse[In: ({type M[x] = Mapping[x, Out]})#M, Out]

Revise Kotlin 1.3 from Scala experience perspective



Thanks to Google Kotlin is producing more and more informational noise. It's interesting to make an assessment spending less than hour to find out where the language is now and comparing it with Scala's experience. I'm not going to review feature by feature - but will try to get impression in offbeat way - choosing the most interesting/unaccustomed feature in 1.3 and trying to assist language's way of development comparing with pseudocode in Scala.

Let's look into upcoming release 1.3.  Kotlin Contracts is looking interesting and probably the biggest KEEP change in the release. Looks like something language specific (haven't seen it in other languages) that improves language itself and may bring some light to the Kotlin's "way of thinking".

Let's run through the proposal to and try to understand what it improves and why.
The first example/motivator is:

Well, it's extremely hard to get why should we write the code like that. There is some background on how Kotlin implements calls of closure blocks (there was similar SIP-21 - SPORES in Scala, but  didn't gain popularity), skipping that - for particular example it feels more natural to use functional approach:
I feel like Kotlin tries to make this code valid:
Kotlin doesn't require us to define val and initialise it immediately, but I don't feel it's nicer and better readable. There is an extra price - method run has a

contract {
    callsInPlace(block, InvocationKind.EXACTLY_ONCE)
}

The next one example is:


As for developer with Scala experience it's hard to understand the problem's domain - but keeping in mind that Kotlin is providing safety for null references problems - there is some sugar for the cases when safety is already checked against the reference and it follows some rules (not null in the example). Looks fine and it's imperative alternative to Monad's approach (will talk about this later). But if it were the production code I would prefer to avoid throwing the exceptions and separate execution of the side effect (println). Same as in previous example there is small complexity via introducing:

Next one example is:

Looks like pattern matching customisation - I had intuitive feeling that it's the way to handle the Union Types but it isn't. As for the guy without commercial experience with Kotlin - all the examples are too artificial to evaluate the syntax sugar coming with contracts. For example the last example looks much better if the pattern matching is applied

Still it could be covered via Option/Either Monads or if there are more possibilities - better to look into Coproduct solutions. All the other examples are following the same paradigm if it were the real product's source code I would prefer to use best practices from Lambda and Categories patterns. Especially that handling nullable (empty) values has the same importance as validation/parallel validation or applying side effects in a good way. There is probably implicit advantage of Kotlin's sugar - allocate less memory and it would be great - but as we will see as soon as we require some feature like ?.let the extra memory allocation is inevitable.



Summary
Overall impression is quite positive. Kotlin is inventing some alternative to Lambda + Category solutions for specific problems.  Nullable reference is really weird solution that appears to be the centrum of all the problems/improvements in the given examples. Fortunately I haven't met NullpointerException quite some time in Scala - but should admit that one of the most used Monads is Option. Another  positive impression is: it's possible to write code in Kotlin from the first minute.

Bonus: Weird nullable types
The Null Safety is looking too noisy but we will play around this area comparing with Scala. Let's imagine we need to read (from property file)  three optional variables and build optional url object based on those. Pretty easy approach checkin all the 3 variables against null and creating the object, Scala allows to write something like:

Or not using the for-comprehensions but Applicative's mapN:

Let's find out can we do something similar. Kotlin doesn't support for comprehensions out of the box. Starting with brute-force solution for 2 params:

Let's use ?.let method:

Looks fine for the case with two params - but wouldn't for more. Let's try to play with Monad's flatMap and unit:

Looks like Maybe/Option Monad:

Looks better - but wouldn't for three and more params. Can we use mapN?

And the usage is:

Works for Pair type only - for more params we need a custom implementation per type, but idea is clean - Kotlin can support functional solutions with Categories. I'm pretty sure there should  be dozen good Functional Programming libraries in Kotlin like Arrow. 


Thursday, 4 October 2018

Protobuf getting metadata at runtime

Sometimes it's important to get the list of fields from Protobuf message.

The good example is Kafka with Protobuf serialiser - when the code generator to Java isn't used - but desc file should be generated anyway. It worse to invest into pact testing of the produced messages - check if the number of sent fields is within the range of defined in proto schema.
Unfortunately Protobuf for JVM doesn't support reading metadata from proto files, but if it's compiled to desc file - it's possible to read it via JVM API - that isn't intuitive:

Friday, 21 September 2018

Down-streaming with Event Sourcing application.

How to apply side effects in Event Sourcing application with desirable/expected consistency level?
It's common and logical to put side effects as a Sink/End of Persistent stream. Usually side effects are dispatched in event-persisted callback handlers. These handlers can also spawn new one sub-streams to react on applied events.

There are different tactics to integrate down-streams to the business flow and this is serious architectural question that makes influence on all the application's design layers.

To review all the possibilities I'm going to use an example of typical event-sourcing application. Let's assume it implements some business utility.  After each event is dispatched we are sending the notification to Business Intelligence services for analyzing and reporting. Imagine the day when you find out - that Business Intelligence operations are fast enough and can produce real-time feedback which can increase the profitability of our business up to X%. We are using BI as an example and placeholder - but in common it can be any downstream service reacting to our business events. Down-stream service are the ones that consume the upstream service.
Just to rephrase it and agree on terminology - I will use BI in all the examples - but it could be any down-stream service.
Given the terminology 'upstream' and 'downstream' it may help to make an analogy with a river. If you drop a message (data) in the river it flows from upstream (initiator) to downstream (receiver).
Here is the typical event-sourcing flow from request to response including other side effects:


Figure 1. Typical event souring flow

After stream emits a command - it should:
  1. Validate input against the current state and business rules then generate mutation events.
  2. Persist an event(s)/state.
  3. Respond to event producer if needed.
  4. Execute other important side-effects
  5. Send events/messages/state to BI (down-stream)
Here is pseudo-code that implements the behavior using akka-persistence:


It's quite logical to apply all side-effects after the mutation has been already persisted (as we operates up-stream and down-stream patterns). In step 3 we send response to the client. Our actor is abstracted from delivering this message, in particular case we rely to "at most once" akka's guarantee. If client expects the different one then we should build business protocol to provide extra guarantees - for example handle timeouts etc.

Here is sequence diagram that helps us to find the flow's operations with different consistency levels:

Figure 2. Consistency per stream flow

The legend for Figure 2 is:
  • A to B is under your API interface consistency - for example it could be Web-sockets protocol with some custom retry policy. We are going to ignore this step flow in examples, ideate it has tolerable consistency and our Architecture design responsible only for the flow from B to E. 
  • C during persistency execution we rely on our persistence journal's consistency - for example Cassandra's.
  • C to D - After journal persisted event - the flow returns to app. "at least once".
  • E - is downstream - it can start the new sub-stream with A to D steps or rely on some other delivery tools, for example Kafka.
As we can see the weakest consistency is between C and D and it determines consistency to BI down-stream. Actually our example implements the first solution with the lowest consistency level but it has own pros as well.
Figure 3. BI side-effects with at most once guarantee

Solution #1: At-most-once, i.e. no guaranteed delivery 

If we send events to BI after the persistence step:
  • The response time from client perspective doesn't depend on down-streaming performance
  • BI delivery failures don't influence on our business flow
  • BI can miss out events (at most once delivery guarantee)
  • Low latency delivery to BI downstream
We can miss out downstream events - for example if we use Kafka producer and rely on Kafka's persistence - when producer.flush operation is failed we might lose the events.

Solution #2: Possible redundant events

Figure 4. Allow redundant events 

If the fact of receiving the events has a priority over dispatching them - it's possible to change the guarantee to "at least once" allowing the duplicates (in case of client retries) and inconsistency between (not all the events are persisted - but all are reported).
This solution provides:
  • Better down-stream but worse response latency (persist waits the flush for downstream completed).
  • At least one guarantee for downstream - but possible inconsistency to state.
  • Can be combined with solution #1 when different events have different guarantees.


Solution #3: Connector to persistence layer

Figure 5. Connector to persistence layer.

Sometimes down-streaming the data is becoming critical for the business flow and more reliable guarantees are wanted. It's possible to use the tools that allow to subscribe for persistence and convert them into stream.

There are plenty of market ready solution.  For example if Cassandra is used as persistence for your journal and BI is a Kafka's consumer - then you can look for different Kafka to Cassandra connectors.  In case of Akka is standing for event-sourcing in your design  - look into Persistence Query

The difference to solution #1 is:
  • Better delivery guarantee to BI (could be tuned to at least once or exact once)
  • Decline in latency
Solution is primely except the cases when latency of BI events is critical. 
You can try to implement custom solution - for example for Cassandra it's possible to read commit log and send the entries into Kafka. It will allow you to tune the latency between Cassandra and Kafka on the lowest possible level - but it will be always a tradeoff - the quicker you make Cassandra flush the data to journal - the slower will be your persistence - but faster stream to BI.

Solution #4 Integrate into persist step

When down-streaming the events is critical part of your flow and not been able to proceed with it means unrecoverable failure for application - then it's logical to amalgamate persist action with down-streaming.


Figure 6. Integrate into persist step


It means we either persist and send events downstream or fail. In case of akka-persistence we can use the journal that persists the events and sends them to Kafka. This combination can be met with DuplicatingJournal and StreamToKafkaJournal.

This solution is:
  • Failing the flow if either persist or send event failures (of course each of them can handle some retry policies etc.)
  • Has a delivery guarantee equaled to chosen downstreaming tool (for Kafka can have "at least once" and "exact once"
  • Response and other side effects might have worse latency than in solution #1 because of persistTime = max(journalPersistDuration, eventSendFlushDuration)
  • Doesn't implement Atomicity guarantees - but explicitly fails in case of inconsistency.

Solution #5 Custom protocol

Let's try to solve the the case when down-streaming consistency and latency are evenly critical - but we can accept eventual consistency in case of failures and quasi real-time for the ninety-nine out of a hundred.
This means we want to continue business flow execution even if down-stream is failing - but we rely on eventual consistency of down-stream - it should recover and continue from the moment it failed as soon as it has been recovered. 

Figure 7. Custom protocol

This is the most expensive solution as it requires implementation of manual delivery protocol to down-stream. Unfortunately Akka's At-Least-Once Delivery is insufficient for our case.
As an example we will use finite-state machine for BI consumer and coupling it's logic to the business flow.

Lets look in details:
  1. Command is emitted.
  2. Each command after validation is checked is it important to be an initialization marker for business flow. Examples of initialization commands are: start of transaction, user creates shopping card, online game round started etc. Initialization event should be delivered with the most desirable guarantee - but it should not make huge impact on main flow latency - because they take only small percentage from all the events.
  3. If event is initial then it has delivery guarantee at least once and business flow just awaits down-stream flush successfully completed. If the event isn't initial then it sent to down-stream in "fire and forget" mode.
  4. Events are persisted.
  5. Side-effects are applied and one of them is sending message to BI. Fire and forget mode like in solution #1.
  6. Sending the response to client.
  7. The most complicated part is BI - it has implements finite state machine that validates the messages been received. 

As an example of that solution we can implement FSM protocol for the service that can adds together integer values and support commands:
  • InitializeTransaction
  • Operand(Int)
  • GetState
Rough algorithms is:
  1. Business flow receives "InitializeTransaction" command
  2. Because we marked this command as initial - it should be sent to BI with at least once consistency.
  3. BI awaits with timeout for the next event - EventPersisted(TransactionInitialized) - in case of timeout it should explicitly fetch the state of transaction from the business flow.
  4. After event is produced by command is persisted - business flow should send the event to BI with fire and forget style (no delivery guarantee).
  5. BI after each event sets timeouts for expecting the next one or final state - it guarantees the eventual consistency between business service and BI down-stream.
  6. BI iterates all the 3 to 5 steps until it gets the final marker - transaction completed.


Solution #6 Kafka streams

Use Kafka streams. Kafka streams became self-sufficient platform for building event sourcing applications and doesn't require hybridizing with other frameworks like akka-persistence.


Summary:

If you don't have to rely on strict delivery guarantees for down-streams use a mix of Solution #1: At-most-once and Solution #2: Possible redundant events. 
If delivery latency isn't critical - but consistency is - use Solution #3: Connector to persistence layer. It should fit to the most of business cases and supported out of the box by the most of event-sourced frameworks. 
Solution #4 Integrate into persist step emulates transaction without rollback possibility and is suitable only for special business cases.
When down-streaming latency is critical Solution #5 Custom protocol or Solution #6 Kafka streams are the ways to go. Implementation and support of custom protocol is the most expensive, additionally it validates your main business flow and theoretically can provide good monitoring feedback. Using the Kafka streams has own limitations and isn't suitable for all the business domains.

Monday, 17 September 2018

ScalaCache - conditional caching

Sometimes it's important to avoid caching some subset from return values. It's easy to implement it using Memoization method memoizeF. When using M[_] like Future - the Failed case won't be cached - and it's possible to convert any value to Failed with special exception marker that contains the value itself and in meantime blocking it from been cached.

Sometimes it's preferable to avoid caching some subset of possible values without deviation to failed case of higher kind wrapper. In case if this condition can be delineated by predicate and you don't want to play with implicit mode: Mode[F] you can mixing small trait to your cache:
This example is based on Caffeine and it isn't caching negative integer values: