Redis Streams and the Unified Log
brandur.org
brandur.org
This. Precisely this. Redis Streams, if done right, will be a killer feature. Redis is already pretty much omnipresent and has a great managed cloud story (e.g. Elasticache). And it's cheap.
We are doing some work at my job right now inspired by this approach (in no small part, that LinkedIn article). I hope we can open-source some of this tech, or at least start blogging about the concepts as nobody really talks about it, and the mechanics are a little tricky to get right.
We currently think the best approach is to bring ideas from DDD (specifically the idea of a bounded context and a domain), a unified log which we call the GEL (for Global Event Log), domain-specific event logs, event-sourcing with CQRS and then some quite lightweight means to produce "projections" (different structural representations of data).
We therefore have a few pieces that get names and become easy to reason about: every event goes to the GEL. Domains have systems that listen to the GEL for events they are interested in and map those events into commands on a command service. Clients can issue commands directly, as well. Commands alter projections and produce events back to the GEL. Clients and domains can query a domain's query service.
It's not really rocket science, these ideas are as old as the oldest mainframe you've ever heard of, but it affords incredible flexibility and unrivalled performance.
We're using Kinesis for our GEL that has a read performance issue, but not a big enough issue for us yet to look at alternatives - the zero-admin cost is attractive to us right now.
For me, now, MVC and CRUD style applications look like really odd anti-patterns. I'm amazed we got as far as we have with them.
I'm sure those patterns are used in great applications but personally I would really think hard if I need the functionality before I went there.
MVC I can understand but how does the unified log pattern/event sourcing help with (or replace) the datastore behind CRUD style applications?
Not gp, but maybe because without book-keeping you lose information in a CRUD style app? The append-only (until/unless you truncate it) log gives you a lot of durability and insight into the system via replaying events from a snapshot + log combination. At least, that's my understanding
Smaller subsets of events are kept in Redshift/Postgres/MongoDB databases, depending on the purpose and how they'll be queried. These keep anywhere from 48 hours to 6 months of data depending on the purpose, sometimes it's just the raw events filtered in a certain way, and sometimes a projection based on those events. Either way, if we catch a bug in our logic we can re-run it for the past year of data or whatever.
We have a public unified log (JSON blobs, behind a CDN) containing all events about package publishes, edits, and deletes. The rest of our endpoints have background jobs reading this log and updating different "views" of the package data. Each event's unique ID is a timestamp of when it was added to the log, much like Redis Streams.
We have found this is a powerful concept which has helped us build a very reliable infrastructure where almost all public endpoints are static JSON blobs served by a CDN. The only compute we need for customer API calls is a search service.
We hope one day that our unified log (called the "catalog") will be used for custom client needs and for package replication (e.g. corporations behind a firewall that can't talk to NuGet.org directly).
API docs: https://docs.microsoft.com/en-us/nuget/api/catalog-resource Catalog root: https://api.nuget.org/v3/catalog0/index.json
"Redis streams aren’t exciting for their innovativeness, but rather than they bring building a unified log architecture within reach of a small and/or inexpensive app. Kafka is infamously difficult to configure and get running, and is expensive to operate once you do. [...] Redis on the other hand is probably already in your stack."
https://github.com/nats-io/nats-streaming-server#nats-stream...
As a producer, you just need to know that your message has definitely been delivered to the log. What consumes it is none of your business.
Have you read https://engineering.linkedin.com/distributed-systems/log-wha... ? It's one of my favourite essays on software architecture of all time, and it goes deep into this concept.
Right I think we're saying the same thing. In the article checkpointing was introduced to enable safely managing the size of the log, i.e. it appears that in this use case _somebody_ needs to know whether the consumers are caught up in order to adhere to at-least-once semantics. It just seemed to me that if you reach the point where you're rolling your own to solve that problem there are other possible solutions.
I haven't read the linked-in piece yet but I intend to. It was also linked in the article. Thanks for linking it here.
All mutations run through the Kafka Streams application. A Relay webapp client communicates with a graphql-java server which resolves GraphQL queries against a gRPC service layer. This service layer either queries Kafka Streams state stores, other state stores, or writes mutations to Kafka, wherein they are processed by the Kafka Streams topology. Another Kafka Streams topology indexes topics to Elasticsearch. Having separate topologies seems to make it easier to reset them independently in case I want to replay state. All of this is glued together with Protocol Buffers, which are a great complement to both Kafka and Elasticsearch.
It is really nice to see Redis join this scene and to see more people getting excited about building applications this way. Looking forward to checking this out.
Pretty neat...
One set of processes retrieves and stores the data in the normalized structure, another set of processes reads and processes them. If there's ever a problem in the processing, the system can fail until someone looks into the error and adds an exception or fixes the processing. This is useful for when you can't be sure the remote side will continue to work as expected (I've found API's to be only slightly more stable than random web scraping. Sure, they may version their APIs, but that doesn't help when they don't communicate new versions or when they are deprecating the old ones and you're hit with a dead API endpoint it out of the blue).
Not only is the processing asynchronous, but the development is so asynchronous, I've added data sources that I knew I would eventually want to process (as adding a data source is fairly easy), only to get around to writing the processor six months later, and the nature of the system means that it's trivial to process the old data. In fact, the normal operation would likely note it is that far out of date and process the old data. The only consideration that needs to be made is whether you want to optimize the back-processing by aggressively caching, since that can greatly speed up the process.
I really think the pattern it is the natural conclusion most people would come to when presented with the right requirements. The real art comes from recognizing when your current problem can map to a pattern such as this when it doesn't appear to initially.
1: e.g. "/$base/$logname/%Y-%m-%d/%H/%Y%m%d%H%M%S_${unique}.json.gz"
And Kafka is hard to deploy ? Please... If you deploy a redis cluster without knowing the pitfalls of distributed systems, you'll have the same problems as with Kafka. Also, soon enough you can deploy brokers over jbods without ZK, what will be the argument then I wonder...
About the use cases where strong guarantees are needed: for instance Disque (http://github.com/antirez/disque) provides strong guarantees of delivery in face of failures, being a pure AP system with synchronous replication of messages. For Redis 4.2 I'm moving Disque as a Redis module. To do this, Redis modules are getting a fully featured Cluster API. This means that it will be possible, for instance, to write a Redis module that orchestrates N Redis masters, using Raft or any other consensus algorithm, as a single distributed system. This will allow to also model strong guarantees easily.
"Pricing for a small Kafka cluster on Heroku costs $100 a month and climbs steeply from there. It’s temping to think you can do it more cheaply yourself, but after factoring in server and personnel costs along with the time it takes to build working expertise in the system, it’ll cost more."
If your project can afford Kafka, use Kafka. This article is about achieving the same pattern in any project that already has Redis.
I think this is largely a rephrasing of one of the findings of dealing with NoSQL vs. Relational DBs, which is that for a subset of problems relational DBs have been used to solve in the past, NoSQLs are perfectly capable of replacing them, but not every problem a relational DB handles can be handled by a NoSQL without investing significant time and effort layering more complex systems on top, at which point you've basically implemented a poorly optimized ad-hoc relational DB. In this case, Redis, with the new stream type is capable of handling a subset of problems that have previously required Kafka to solve, but it isn't itself a replacement for Kafka in all situations since there are several important features of Kafka that aren't available in Redis.
The staged log record idea is interesting but it makes the database bottleneck bigger specially if a lot of different transaction (spread across many tables) now is bottlenecked on the same stagedlog table.
Maybe I'm missing something here?
The other big issue is that even if the insert into the shared table is fast (simpler structure) the transaction hence the commit can be long due to the complexity of the application transaction.
I will be happy to be corrected though if we are doing something wrong.
What sort of thing are you after? Like a getting started for production without creating footguns style how too?
Perhaps you had a different issue though?
Doesn't this constraint only hold true if the XADDs to redis are completed synchronously?
Can someone translate the example "rocket-rides-unified" code for non-Ruby coders can read it too? (preferable to a C-like syntax language like Node.js, PHP, Go or Java)