17 ms·
Kafka Is Not a Database
- fouc 6y agoAny sufficiently complex software will end up implementing a database.
- dgb23 6y agoThe article links to this talk[0], which has a funny and interesting sounding title: "Did you accidentally build a database?" [0] https://www.oreilly.com/library/view/strata-hadoop/9781491944608/video244677.html https://www.oreilly.com/library/view/strata-hadoop/978149194...
- revertts 6y agoThat link's a 3min clip for non-subscribers, but the full talk is here https://www.youtube.com/watch?v=Bz2EXg0Fy98 https://www.youtube.com/watch?v=Bz2EXg0Fy98
- quickthrower2 6y agoI’d qualify that as “complex infrastructure software”
- joking 6y agoneither has to be.
- tacitusarc 6y agoI think because software engineers tend to excel at pattern recognition, oftentimes solutions to different problems appear so similar that it seems like with a small amount of abstraction, they can be reused. But it's a trap! Everything abstracted to the highest level is the same, but problems aren't solved at the highest level. The devil, as they say, is in the details.
- UK-Al05 6y agoThis is lack of abstraction. You can certainly fix this in kafka using various hacks, but its implementation of an abstraction you can get in a standard db for free. Funnily enough a list of events is pretty much what a transaction log is in a standard db. Although the events have more of a business meaning. In many ways event sourcing is removing a lot of abstraction databases give you.
- mrkeen 6y ago> a standard db Yes, ACID works for one database. Many databases? Not so much. > Funnily enough a list of events is pretty much what a transaction log is in a standard db. When I "SELECT Balance WHERE user = 12345", I usually just get back a balance, I don't get back the transaction log. If nothing else, adopting the Kafka model gets your teammates to append updates to your ledger, rather than changing values in-place.
- UK-Al05 6y agoYes that my point. You get half of a database. The transaction log. Not the up to date state.
- UK-Al05 6y agoOne way around this is to make sure your kafka command streams are processed in order, in serial partitioned by an id where you want the concurrency control. Normally you only want concurrency control within certain boundaries. By figuring out the minimum amount transaction and concurrency boundaries you can inch out quite a bit of performance.
- loopz 6y agoSure, but that defeats the quest for horizontal scalability. You can build highly performant systems based on serial execution, but not sure this is an area where Kafka excels particularly.
- UK-Al05 6y agoThat's why you partition by some id. Say stock SKU id for stock control. Then you can handle other SKUs in parallel. It's only in serial for a single SKU. That's probably the maximum performance potential your going to get in a traditional db anyway.
- edbrown23 6y agoThis definitely seems like the "Kafka" way to solve this problem, but I fear there are implications to this partitioning scheme I'd love to see answered. For example, partition counts aren't infinite, and aren't easily adjusted after the fact. So if you choose, say, 10 partitions originally, for a SKU space that is nearly infinite, then in reality you can only handle 10 parallel streams of work. Any SKU that is partitioned behind a bit of slow work is then blocked by that work. It's doable to repartition to 100 partitions or more, but you basically need to replay the work kept in the log based on 10 partitions onto the new 100 partitions, and that operation gets more expensive over time. Then of course you're basically stuck again once your traffic increases to a high enough level that the original problem returns. If the unit of horizontal scaling is the partition, but the partition count can't be easily changed, consumers eventually lose their horizontal scalability in Kafka, from my perspective.
- jkarneges 6y agoAnother potential misuse of Kafka I've been wondering about is how a single Kafka instance/cluster is often shared by multiple microservices. On one hand the ability to connect multiple microservices to a central message broker is convenient, but on the the other hand this goes against the microservice philosophy of not sharing subcomponents (databases, etc). I wonder where the lines should be drawn.
- ivalm 6y agoWait, what? Isn’t the whole point of having multiple publishers/subscribers?
- mrweasel 6y agoI think the point was about using a single cluster for multiple topics, for different services. Depending on the scenario I can see the point. If the micro services are all part of the larger overall solution, having a single cluster is perfectly fine. Using the same cluster for multiple "product" is a little like having one central database server for a number of different solutions. You can do it, but it potentially become a bottleneck or a central point for your different solutions to impact performance of each other.
- jkarneges 6y agoI'd agree there is an arguable difference between sharing a server vs sharing data within the server. Bottleneck issues aside, letting two microservices connect to the same Postgres cluster but access different "databases" (collection of tables) within that cluster could be considered an acceptable data separation. Certainly with multi-tenant DBaaS systems there may be some server sharing by unrelated microservices/customers. Whereas letting two microservices access the same database tables would probably be frowned upon. Nevertheless, sharing the same Kafka topics between microservices seems to be a common thing to do.
- ivalm 6y ago> Whereas letting two microservices access the same database tables would probably be frowned upon. > Nevertheless, sharing the same Kafka topics between microservices seems to be a common thing to do. I think if it is part of one whole isn’t this fine? You have one service that generates customer facing output, you may have another service that powers analytics/dashboards you may have yet another service that ETLs data into some data mart. Why wouldn’t they touch the same table/subscribe to the same topic (since they just need read-only access to the data)? Genuinely curious what the problem is except for bottleneck/performance; and if it just bottleneck then wouldn’t scaling horizontally solve it?
- tutfbhuf 6y agoWell, then you have never heard of ksqlDB. It adds SQL and DB features to Kafka. It is backed by Confluent (LinkedIn) same company that developed Kafka initially. https://ksqldb.io https://ksqldb.io
- albertwang 6y agoWe’re very aware of ksqlDB. I would recommend this video from last week where Matthias talks about some of the strengths and weaknesses of ksql: https://youtu.be/KUQuegJ4do8 https://youtu.be/KUQuegJ4do8
- cirego 6y agoMy understanding is that ksqlDB is a read-only interface on top of streams and only helps people write better consumers. The problems mentioned in the blog post relate to producers.
- Kalium 6y agoAs recently as last year, I worked for a company where the Chief Architect, in his infinite wisdom, had decided that a database was a silly legacy thing. The future looked like Kafka streams, with each service being a function against Kafka streams, and data retention set to infinite. Predictably, this setup ran into an interesting assortment of issues. There were no real transactions, no ensured consistency, and no referential integrity. There was also no authentication or authorization, because a default-configured deployment of Kafka from Confluent happily neglects such trivial details. To say this was a vast mess would be to put it lightly. It was a nightmare to code against once you left the fantasy world of functional programming nirvana and encountered real requirements. It meant pushing a whole series of concerns that isolation addresses into application code... or not addressing them at all. Teams routinely relied on one another's internal kafka streams. It was a GDPR nightmare. Kafka Connect was deployed to bridge between Kafka and some real databases. This was its own mess. Kafka, I have learned, is a very powerful tool. And like all shiny new tools, deeply prone to misuse.
- JMTQp8lwXL 6y agoInstead of single architects, companies need architect boards. And they can vote on these ideas before a single individual becomes a single point of failure. Expecting 1 person to make 100% correct decisions all the time is too much expectation for one person. People go down rabbit holes and they have weird takeaways, like replace all the databases with queues.
- Kalium 6y agoI agree in abstract, but in practice it's quite difficult to set up a successful democratic architecture board. You need teams or departments that all have architects, and an engineering organization where both managers and engineers accept a degree of centralized technical leadership. Getting there is, in my opinion, the work of years. It's especially challenging because spinning up such a board requires a person who can run it single-handedly. In this particular company, the Chief Architect in theory had a group around him. They did nothing to check his poor decisions, and from the outside seemed primarily interested in living in the functional programming nirvana he promised them.
- zaphar 6y agoIf it stores data it's a database. Filesystems are databases, MongoDB is a database. LevelDB is a database. Postgres and MySQL are databases. Kafka is a database. They are all very different in features and functionality though. What the authors mean is that kafka is not a traditional database and doesn't solve the same problems that traditional databases solve. Which is a useful distinction to make but is not the distinction they make. The reality is that database is now a very general term and for many usecases you can choose to special purpose databases for what you need.
- yumaikas 6y agoI think I'd differentiate between a database and a data store. I'd argue that a filesystem is a data store, rather than a database.
- edgyquant 6y agoI’d say initially file systems were data stores but once they developed hierarchies they became more akin to a database. I’m not sure there’s a huge difference but it seems a database is a collection of data stores (though there or probably a more technical and correct definition.)
- Hamuko 6y agoSome say that Microsoft Excel is the world's most popular database engine.
- Baeocystin 6y agoHonestly, I'd agree with that, too, at least in a colloquial sense.
- victor106 6y agoThe most popular database where the number of concurrent users = 1
- 6y ago
- detay 6y agousing the any tool for correct problem requires skills.
- diehunde 6y agoRelevant to the discussion: Martin Kleppmann | Kafka Summit SF 2018 Keynote (Is Kafka a Database?) [1] [1] https://www.youtube.com/watch?v=v2RJQELoM6Y https://www.youtube.com/watch?v=v2RJQELoM6Y
- fredliu 6y agoKafka is essentially commit logs, which are at the core of any traditional database engines. Streaming is just turning the gut of DB inside out (mostly for scalability reasons), while DB is wrapped up commit logs that provides higher level functionalities (ACID, Transactions, etc.). It's two sides of the same coin, yin and yang of the same thing... But on the practical side of things, yes, if what you needed more are indeed what's described in this article, your life would be easier with a traditional DB.
- j-pb 6y agoTbh, It's a weird blog post coming from the materialize folks, considering they know better. The "event sourced" arch they sketched is missing pieces. Normaly you'd have single writer instances that are locked to the corresponding kafka partition, which ensure strong transactional guarantees, IF you need them. Throwing shade for maketings sake is something that they should be above. I mean c'mon, I'd argue that Postgres enhanced with Materialize isn't a database anymore either, but in a good sense! It's building material. A hybrid between MQ, DB, backend logic & frontend logic. The reduction in application logic and the increase in reliability you can get from reactive systems is insane. SQL is declarative, reactive Materialize streams are declarative on a whole new level. Once that tech makes it into other parts of computing like the frontend, development will be so much better, less code, less bug, a lot more fun. Imagine that your react component could simply declare all the data it needs from a db, and the system will figure out all the caching and rerendering. So yeah, they have awesome tech with many advantages, so I don't get why they bad-mouth other architectures.
- dkhenry 6y agoI am not terribly surprised. The materialize team was previously at CockroachDB which also had a habit of putting out marketing material like this.
- j-pb 6y agoMaybe I'm biased because I'm such a huge Frank McSherry fanboy. The differential dataflow work he does in rust is simply awesome, and he writes great papers too! Little known fun fact, the rust type- and borrow-checker uses a datalog engine internally to express typing rules, and that engine was written and improved by Frank McSherry. So whenever you hit compile on a rust program, you're using a tiny bit of Materialize tech.
- _jal 6y agoI ignored them for a long time because of it. I just assumed they were lightweights being annoying because they didn't have anything else. It serves as anti-marketing, at least to me.
- 6y ago
- Cojen 6y agoJim Gray disagrees: https://arxiv.org/ftp/cs/papers/0701/0701158.pdf https://arxiv.org/ftp/cs/papers/0701/0701158.pdf
- georgewfraser 6y agoIt's funny that you use that example, we actually cited that in an earlier draft of this post. Despite the seemingly opposite title "Queues are Databases", that note actually makes many of the same arguments, that message brokers are missing much of the functionality of database management systems and this is a problem.
- based2 6y agohttps://www.postgresql.org/docs/10/rules-materializedviews.html https://www.postgresql.org/docs/10/rules-materializedviews.h...
- Spivak 6y agoI feel like the inventory thing is a bit of a straw-man because the situation is set out in such a way that you need transactions for it to work. If you find yourself wishing you had a global write-lock on a topic to then of course it won't work. Modeling your data for Kafka is work just the same as it is for MySQL. Of course it might not be the best tool for the job but you should at least give it a fair shake. You should be able to post "buy" messages to a topic without fear that it messes up your data integrity. Who cares if two people are fighting over the last item? You have a durable log. Post both "buys" and wait for the "confirm" message from a consumer that's reading the log at that point in time, validates, and confirms or rejects the buys. At the point that the buy reaches a consumer there is enough information to know for sure whether it's valid or not. Both of the buy events happened and should be recorded whether they can be fulfilled or not.
- jacques_chester 6y ago> Who cares if two people are fighting over the last item? The two people, at least. Customers tend to be a bit underwhelmed by "well, the CAP theorem..." as a customer service script.
- Spivak 6y agoNone of this shows up as user-facing any differently than a relational database. No CAP theorem at all. Kafka: User clicks buy and it shows “processing” which behind the scenes posts the buy message and waits for a “confirmed” message. When it’s confirmed user is directed to success! If someone else posts the buy before them they get back a “failed: sold out” message. Relational: User clicks buy and it shows “processing” which behind the scenes tries to get a lock on the db, looks at inventory, updates it if there’s still one available, and creates a row in the purchases table. If all this works the user is directed to success. If by the time the lock was acquired the inventory was zero the server returns “failure: sold out”.
- jacques_chester 6y agoThe CAP theorem line was smart-arsery. The thing here is that the database can update the cart and the inventory in one logical step, to the exclusion of others. The Kafka approach doesn't guarantee that out of the box, leading to the creation of de facto locking protocols (write cart intent, read cart intent, write inventory intent ...). A traditional database does that for you with selectable levels of guarantees.
- AndrewKemendo 6y agoAlternatively from Jay Krebs [1] a much more thorough and nuanced discussion that is probably the best send-up on this topic. "So is it crazy to do this? The answer is no, there’s nothing crazy about storing data in Kafka: it works well for this because it was designed to do it. Data in Kafka is persisted to disk, checksummed, and replicated for fault tolerance. Accumulating more stored data doesn’t make it slower. There are Kafka clusters running in production with over a petabyte of stored data." [1] https://www.confluent.io/blog/okay-store-data-apache-kafka/ https://www.confluent.io/blog/okay-store-data-apache-kafka/
- halbritt 6y agoI recommend his book, "I heart logs". It's a short read, but changed my perspective.
- pram 6y agoLog compaction definitely isn’t problem free. I’d say it was crazy when he wrote that. We ran a cluster with lots of compacted topics, hundreds of terabytes of data. At the time it would make broker startup insanely slow. An unclean startup could literally take an hour to go through all the compacted partitions. It was awful.
- georgewfraser 6y agoThat post explains that there are scenarios where it makes sense to store data permanently in Kafka. "Kafka is Not a Database" makes a different point, which is that Kafka doesn't solve any of the hard transaction-processing problems that database systems do, so it's not an alternative to a traditional DBMS. This is not a straw man---Confluent advocates for "turning the database inside out" all over their marketing materials and conference talks.
- cfontes 6y agoIt actually does with |Exactly-once Semantics| in fact I've been using as the single source of truth in a cash management system for almost 2 years without a single issue related to transactions.
- je42 6y agoOk. I admit using Kafka as DB is not straight forward but just stating it doesn't provide ACID functionality is not enough. The example they give is very simplistic. With the correct design of kafka topics and events the problem of the example can be fixed. And according to oracle https://www.oracle.com/database/what-is-database/ https://www.oracle.com/database/what-is-database/ : > A database is an organized collection of structured information, or data, typically stored electronically in a computer system. So Kafka clearly fits that definition.
- halbritt 6y agoThe author presumes that every use case requires a transactional database. ACID is nice, especially if it's needed, but generally not needed, especially in many streaming data applications for which Kafka is most suitable.
- satisfaction 6y agoI'm guilty of using it as a DB, my home weather station writes to kafka topics, sometimes the postgres instance is down for months, no problems letting the kafka topic store the data until i get around to rebooting pg and restarting the connector.
- dathinab 6y agoDidn't some newspaper use Kafka to store the newspapers they released in order or something similar? (I think it was the New York Times, maybe??). Honestly as long as you don't use it as a general purpose database, it might very well be the best choice for your use-case.
- je42 6y agoIndeed. Known the tools and apply the one that has least cons and most pros ;). Also, a good starter to for knowing how to use a Streaming Store like Kafka as DB is the video Database inside-out. https://www.youtube.com/watch?v=fU9hR3kiOK0 https://www.youtube.com/watch?v=fU9hR3kiOK0
- barnabask 6y ago
- megachucky 6y agoHere is my answer if Kafka is a database: (the answer is yes, but you should still not try to replace every other database) https://www.kai-waehner.de/blog/2020/03/12/can-apache-kafka-replace-database-acid-storage-transactions-sql-nosql-data-lake/ https://www.kai-waehner.de/blog/2020/03/12/can-apache-kafka-... Would be curious what the people here think about my post, too. Disclaimer: I work for Confluent - and I am also happy about critical feedback.
- hasanic 6y agoI mean, duh? Does Apache Kafka ever made the claim that it is a database? Other things that are not a database: Apache Traffic Server, Apache Mahout, Apache Jakarta, Apache ActiveMQ... hundreds of these exist.
- dang 6y agoThe title sounds generic, but the article makes it clear that it's responding to specific proposed uses of Kafka.
- hodgesrm 6y agoIs this really a thing? Do people really try to use Kafka as the system of record for financial transactions or similar data?
- mrkeen 6y agoHell yes. Best thing I ever did. Updating balances using an RDBMS is like managing your finances with pencils and erasers. Unless you somehow ban the UPDATE statement. Updating balances with Kafka is like working in pen. You can't[1] change the ledger lines, you can only add corrections after the fact. [1] Yes, Kafka records can be updated/deleted depending on configuration. But when your codebase is written around append-only operations, in-place mutations are hard and corrections are easy, so your fellow programmers fall into the 'pit of success'.
- hodgesrm 6y agoI would just put the ledger in a database table if it's that important and maintain the current state of the account in a separate table. ACID transactions and database constraints make this kind of consistency easier to achieve than many alternatives. It's also easier to prove correctness since you can run queries that return consistent results thanks to the isolation guaranteed by ACID. (Modulo some corner cases that are not hard to work around.) Just my $0.02.
- mrkeen 6y ago> ACID transactions and database constraints make this kind of consistency easier to achieve than many alternatives. If your company only runs one database. > It's also easier to prove correctness RDBMSs are wonderful and I don't consider them unreliable at all. But I can't prove the correctness of my teammate's actions. I want them to show their working by putting their updates onto an append-only ledger.
- jacques_chester 6y ago> If your company only runs one database. I think the same argument can be made with "only one Kafka cluster" and "only one blockchain".
- amai 6y agoElasticsearch is also not a database.
- charlieflowers 6y agoCould you elaborate on that? Because it sure seems like a DB to me. Not a relational one. And should not replace a typical CRUD OLTP DB. But it sure seems like a no-sql DB to me.
- amai 6y agoElasticsearch is a search engine. It is optimised to return the best fitting results to a query at the first page. But if you want to retrieve all results for a query (which is a common usecase for DBs) the performance of elasticsearch will plummet. For example: Try to retrieve a list of all the 10000 movies you stored in an Elasticsearch index. You will get the first 100 results easily, but is you scroll through the results, you will notice that elasticsearch will become very slow. It is not optimized for that use case. Otherwise it would be called ElasticDB.
- deleted 6y ago[deleted]
- jgraettinger1 6y agoThis post doesn't mention the _actual_ answer, which is to: 1) Write a event recording a _desire_ to checkout. 2) Build a view of checkout decisions, which compares requests against inventory levels and produces checkout _results_. This is a stateful stream/stream join. 3) Read out the checkout decision to respond to the user, or send them an email, or whatever. CDC is great and all, too, but there are architectures where ^ makes more sense than sticking a database in front. Admittedly working up highly available, stateful stream-stream joins which aren't challenging to operate in production is... hard, but getting better.
- fuckadtech 6y agoHard is an understatement. Particularly so if you are using Kafka Streams to attempt to run a highly available, fault tolerant, zero downtime, etc., service. The race condition and compaction bugs in that library are not fun to debug.
- soumyadeb 6y agoThe architecture of dumping events into Kafka and creating materialized views is a perfect choice for many use cases - e.g. collecting clickstream data and building analytical reports. If ACID is a prerequisite, then lot of things won't classify as databases - None of Mongo, Cassandra, ElasticSearch etc. Not even many data-warehouses.
- EamonnMR 6y agoKafka is a very nice communication channel. You can dump the results into a database and query it if you need a database.
- arthurcolle 6y agoI had an issue with RabbitMQ where I didn't know how my consumer was going to use the data that I was writing to a queue yet (from a producer that was listening on a SocketIO or WebSockets stream), and I was kind of just going to figure it out in an hour or something. Eventually, my buffer ran out of memory and I couldn't write anything else to it, and it was dropping lots of messages. I was bummed. Is there a way to avoid this in Kafka?
- atmosx 6y agoRabbitMQ and Kafka serve different needs. RabbitMQ is the most widely used open source implementation of the AMQP protocol. It is slower but can support complex routing scenarios internally and handle situations were at-least-once-delivery guarantees are important. RabbitMQ supports on-disk persistent queues, which you can tune if you like. Compared to Kafka, RabbitMQ is slow in terms of volume that can be managed per queue. Kafka is fast because it is horizontally scalable and you have parallel producers and consumers per topic. You can tune the speed and move needle where you need between consistency and availability. However, if you want things like at-least-once-delivery and such, you'll have to use the building blocks kafka gives you, but ultimately you'll have to handle this on the application side. Regarding storage, by default kafka stores data for 7 days. IIRC the NY Times stores articles from 1970 onwards on kafka clusters. The storage is horizontally scalable and durable. This is a common use case. As many have pointed out, the cluster setup depends highly on you needs. We store data for 7 days in kafka as well and it's in the order of 500GB or more per node. Looks like you have a configuration issue. You can configure rabbitMQ to store queues on the hard disk and with a quick calculation you can make sure you have enough space for 10 or 150 hours of data. I don't see any reason to switch to kafka, a different tool with different characteristics, just because you need more storage.
- shay_ker 6y agoThis is maybe a silly question, but what's the difference between the timely dataflow that Materialize uses and Spark's execution engine? From my understanding they're doing very similar things - break down a sequence of functions on a stream of data, parallelize them on several machines, and then gather the results. I understand that the feature set of timely dataflow is more flexible than Spark - I just don't understand why (I couldn't figure it out from the paper, academic papers really go over my head).
- mamon 6y agoThere's no difference really. All "Big Data" (tm) tools are trying to capitalize on the hype, so Kafka adds database capabilities, while Spark adds Streaming. At some point they will reach feature parity.
- frankmcsherry 6y agoThere are a few differences, the main one between Spark and timely dataflow is that TD operators can be stateful, and so can respond to new rounds of input data in time proportional to the new input data, rather than that plus accumulated state. So, streaming one new record in and seeing how this changes the results of a multi-way join with many other large relations can happen in milliseconds in TD, vs batch systems which will re-read the large inputs as well. This isn't a fundamentally new difference; Flink had this difference from Spark as far back as 2014. There are other differences between Flink and TD that have to do with state sharing and iteration, but I'd crack open the papers and check out the obligatory "related work" sections each should have. For example, here's the first para of the Related Work section from the Naiad paper: > Dataflow Recent systems such as CIEL [30], Spark [42], Spark Streaming [43], and Optimus [19] extend acyclic batch dataflow [15, 18] to allow dynamic modification of the dataflow graph, and thus support iteration and incremental computation without adding cycles to the dataflow. By adopting a batch-computation model, these systems inherit powerful existing techniques including fault tolerance with parallel recovery; in exchange each requires centralized modifications to the dataflow graph, which introduce substantial overhead that Naiad avoids. For example, Spark Streaming can process incremental updates in around one second, while in Section 6 we show that Naiad can iterate and perform incremental updates in tens of milliseconds.
- theptip 6y agoThis is a bit dumbed down, and ignores the domain terminology required to properly discuss the trade-offs here (which is puzzling given that it links to a post by Aphyr, where you can find incredibly thorough discussions around isolation levels and anomalies). > The fundamental problem with using Kafka as your primary data store is it provides no isolation. This is false. I can only assume the author doesn't know about the Kafka transactions feature? To be specific, Kafka's transaction machinery offers read-committed isolation, and you get read-uncommitted by default if you don't opt-in to use that transaction machinery (the docs: https://kafka.apache.org/0110/javadoc/index.html?org/apache/kafka/clients/consumer/KafkaConsumer.html https://kafka.apache.org/0110/javadoc/index.html?org/apache/...). Depending on your workload, read-committed might be sufficient for correctness, in which case you can absolutely use Kafka as your database. Of course, proving that your application is sound with just read-committed isolation is can be challenging, not to mention testing that your application continues to be sound as new features are added. Because of that, in general I think that the underlying point of this article is probably correct, in that you probably shouldn't use Kafka as your database -- but for certain applications / use-cases it's a completely valid system design choice. More generally this is an area that many applications get wrong by using the wrong isolation levels, because most frameworks encourage incorrect implementations by their unsafe defaults; e.g. see the classic "Feral concurrency control" paper http://www.bailis.org/papers/feral-sigmod2015.pdf http://www.bailis.org/papers/feral-sigmod2015.pdf. So I think the general message of "don't use Kafka as your DB unless you know enough about consistency to convince yourself that read-committed isolation is and will always be sufficient for your usecase" would be more appropriate (though it's certainly a less snappy title).
- georgewfraser 6y ago"Read-committed isolation" is not a meaningful implementation of transactions. If you can't do read, then a write, while guaranteeing the database didn't change in between, then you don't really have transactions.
- deleted 6y ago[deleted]
- 6y ago
- jdmichal 6y agoSo, the problem really being addressed but not named is that eventing systems give eventual consistency. But sometimes that's not good enough. And it's OK to admit that and bring in another technology when you need a stronger guarantee than that. The example I was taught with was a booking system, where the inventory management system-of-record was separate from the search system. Search does not need 100% up-to-date inventory. A delay between the last item being booked and it being removed from the search results is acceptable. In fact, it has to be acceptable, because it can happen anyway. If someone books the last item after another hit the search button... There's nothing the system can do about that. When actually committing a booking, however, then that must be atomically done within the inventory management system. So, to bring it home, it's OK for the search system to be eventually consistent against bookings, and read bookings off of an event stream to update its internal tracking. However, the bookings themselves cannot be eventually consistent without risking a double-booking.
- vladsanchez 6y agoI want to Upvote this more than once. So much facts into a condensed into a small essay. Good job! Money quote: "Event-sourced architectures like these suffer many such isolation anomalies, which constantly gaslight users with “time travel” behavior that we’re all familiar with."
- rowland66 6y agoMaybe I am just old because I had to Google what gaslighting meant, but as best I can tell getting gaslighted by your system architecture is really stretching the meaning of the term gaslighting to the point of absurdity.
- vladsanchez 6y agoI concur, but we both got it's a metaphor. ;)
- somurzakov 6y agoHow does using Kafka come into play when implementing Actor model? Anyone successfully implemented actor model framework over kafka? interested in learning others' experience
- lmm 6y ago> The problem we now have is called write skew. Our reads from the inventory view can be out of date by the time the checkout event is processed. If two users try to buy the same item at nearly the same time, they will both succeed, and we won’t have enough inventory for them both. And you'll have exactly the same problem if you're using a traditional ACID database: the user saw the item as being available, clicked buy, but it was unavailable by the they went to get it. Using an ACID database doesn't gain you anything; you might as well just use Kafka for everything.
- ben509 6y agoThe ACID system can guarantee nobody is billed for an item that you can't deliver. If you want a more user-friendly guarantee, you can reserve it when it's added to the cart.
- lmm 6y ago> If you want a more user-friendly guarantee, you can reserve it when it's added to the cart. If you open an ACID transaction when the user adds something to the cart and don't close it until they check out, you'll find your database gets locked up pretty quickly. So you can't actually use the ACID transactions to implement the behaviour you want - you have to implement some kind of reserve/commit semantics in userspace, whether you're using an ACID database or not.
- ben509 6y agoYou don't need to hold the transaction open the entire time, you just need the inventory count to be correct. You track the inventory and reservations. Taking a reservation checks that inventory is available. With row-level locking, only that inventory item is locked. If that fails, it can search for timed out reservations, update the inventory and try again. If it succeeds, it decrements the inventory then adds a reservation. At that point, the transaction can close, and in the common case, you only held locks long enough to update a row in inventory and add a row to reservations.