It’s Okay to Store Data in Apache Kafka (2017)
confluent.io
confluent.io
Event sourcing is useful, but using it as a source of truth data store in itself instead of e.g. an occasional journalling mechanism seems pretty fraught.
Wherever you are on that spectrum, there is code. Exercising it may slow down your test suite. It might need to be touched when you update lint rules (at least to add annotations to shut off the new rules).
Code nearly always represents some amount of debt. Code to deal with history, more so.
When a schema changes for an all, someone has to migrate it and the data, unless the design is that all incompatible data is dropped. Either that’s the codebase itself or the dev team is punting to the user. Punting to the user just shifts the burden around.
If the R&D team shipped a shitty / incomplete schema to get the software out the door that they need to change later, then yes, that is technical debt - something they’ll eventually need to repay.
If requirements evolve over time and thus the schema needs to, that is not necessarily technical debt in the usual sense, which usually implies a temporary technical compromise for the purposes of expediency / getting something out the door.
I suppose I agree there is a trade off here and a penalty - after many years the migrations get slow to apply etc, and you could say that checkpointing the schema every few versions and preventing upgrades from any prior point is a way of cleaning up the “migration debt”.
But people like to suggest that there’s some other way , ie. the OP saying “ If you get your domain model wrong at the start, it will haunt you forever.”.... I have never seen any system get the domain model perfectly right at the start!
Financially speaking, it's more like an annuity than a credit card.
You highlight the "need to preserve the conversion code for as long as the history exists". This is a point so decisive that I even wonder if Kafka should provide some support to keep that relationship.
Yes, I believe with event sourcing you typically do exactly that. The very point of it is to use the event log as the source of truth, not just an “occasional journaling mechanism.”
Fortunately, Kafka's already thought about this, so, when the time comes, it should be pretty easy to migrate to an architecture where you're no longer using it as your long-term source of truth.
I came across it in the .net world with cqrs + event sourcing combo long before Kafka became popular. People put a lot effort in to move away traditional current state storage. Often used in industries that have a lot of regulation.
Otherwise it's not different from a audit log.
In that situation, how do you handle migrations on that database, or building a new copy of it from scratch via the log? Your log will have historical data in a different shape than the schema expects, so you'll end up in an uncomfortable "replay a little, migrate, repeat" situation when reading historical data into the index for any reason.
> If Kafka isn’t going to be a universal format for queries, what is it? I think it makes more sense to think of your datacenter as a giant database, in that database Kafka is the commit log, and these various storage systems are kinds of derived indexes or views. It’s true that a log like Kafka can be fantastic primitive for building a database, but the queries are still served out of indexes that were built for the appropriate access pattern.
A wonderful insight!
This is always going to be true. Either your system is built in a way that even the most ignorant of users cannot do permanent damage, or you design it to be understood and run by a small group of people who fully comprehend what they're doing. Of course neither option is easy, but the alternative is a prolific generator of daily trouble.
According to https://issues.apache.org/jira/browse/CASSANDRA-8844: CDC seems a lot less capable than Kafka:
> Cassandra would only write to the CDC log, and never delete from it. > Cleaning up consumed logfiles would be the client daemon's responibility
on the other hand there is the Cassandra CDC Kafka Connector, that seems to do some work...
in the i still don't see why i would want to use Cassandra over Kafka for storing data, when i care about durability and streaming of append only data.
Postgres, Spanner, etc. come to mind as options if you need to store critical append-only data and also have it be queryable in real-time.
Or use a durable unbounded solution like Pravega instead of Kafka? https://news.ycombinator.com/item?id=20059006#20062821
edit: this could be problematic especially when data privacy laws are involved. See the comments from gopalv below.
You should provide examples of either better or cheaper to substantiate your claim.
treating your message bus as a infinite storage system is going give you a bad time.
The data retention policies for users who leave the platform is set to 30 days by GDPR, which requires the data to be deleted and expunged by the storage system in ways in which normal users cannot recover - or in another way, the data needs to "offline".
This is not actually that "all data older than 30 days is thrown away", but that "all data for users who have deleted their accounts need to be deleted from the start of their account creation, 30 days after they say forget-me".
Kafka (or Pulsar or Pravega) or any of the other immutable commit log implementations make forgetting a small slice of data from a large set a complicated and nearly impossible task to accomplish.
You can accomplish some part of this with log compaction assuming you have a definite primary key for all updates (i.e the log needs to be partitioned on a key to do compaction along that key). If there's a way I could declare a primary key as a device+event_timestamp+metric, but delete by a user_id column, let me know and I'll be happy to find out how to do it.
However in its original form, Kafka is still very useful.
Being able to hold data in Kafka in those periods of time is extremely valuable and naturally lets the system lose part of its state with the ability to replay itself back into the same state from a saved checkpoint.
If you store 7+ days of Kafka data and flush the newly arrived data-set into a persistent, but mutable columnar store every day & maintain the partition/offsets on commit, then you can recover from a complete loss of the mutable store's in-memory data by replaying the log from where you left off.
The row-major nature of its storage still hurts though if you plan to do all your analysis off it directly, because you'll burn through the disk bandwidth for no good reason.
We use Kafka as our storage for almost everything, and we managed to solve this by encrypting all user data that is relevant to GDPR and trowing away the key when asked for a removal.
if a user asks to be forgotten, we commit a empty privacy key for this user and compress the privacykeys topic and all is done, no service will be able to decrypt it anymore.
So far it has been a good solution and it was easy to implement on all our services.
First, it requires some policing of Kafka use. It's easy for developers to slip up and some PII to spill into the append-only data systems.
Second, your developers will have to handle for what happens when the key is deleted. The happy-path of fetching data, fetching key, and applying will fail quite hard the first time the rare event of a key deletion comes around.
According to France's GDPR supervisory authority, CNIL, organisations don't have to delete backups when complying with the right to erasure. Nonetheless, they must clearly explain to the data subject that backups will be kept for a specified length of time (outlined in your retention policy).
That quickly gets unwieldy and I discard that first instinct. But now people are doing it for real!
That said, I do find pretty quickly you need to get your data out of Kafka in order to query it. Keeping your data in Kafka forever and re-upserting it into your database has mileage though.
Those are tractable but hard to solve; log compaction is not a silver bullet and unless you think really hard about how your data changes over time, you may end up storing more of it than you expect if you use the log as an eternal source of truth--compaction or not.