12 ms·
Issues we've encountered while building a Kafka based data processing pipeline
- jpgvm 5y agoWhat you want is called Apache Pulsar. Log when you need it, work queue when you need that instead. As for ensuring transactional consistency that isn't so bad, you can use an table to track offset inserts making sure you verify from that before you update consumer offsets (or Pulsar subscription if you go that route).
- HelloNurse 5y agoThe only time I used Kafka, it was involuntarily (included for the sake of fashion in some complicated IBM product, where it hid among WebSphere, DB2 and other bigger elephants) and it ran my server out of disk space because due to a bug ridiculously massive temporary files weren't erased. Needless to say, I wasn't impressed: just one more hazard to worry about.
- kitd 5y agodue to a bug Data retention time is Kafka config 101. Are you sure it was a bug?
- geodel 5y agoConsidering how half-assed Kafka is in general, that it needs all clients code changes when Kafka servers are upgraded. It is very likely that user hit Kafka bug.
- kitd 5y agoCitation needed. New server versions are protocol backwards compatible so I'm not sure what you're referring to. Ofc, if you downgraded a server without changing the client, that may cause problems, but tbh that's hardly Kafka's fault.
- deleted 5y ago[deleted]
- EdwardDiego 5y ago> it needs all clients code changes when Kafka servers are upgraded It absolutely doesn't. Message formats predating Kafka 0.11 are only just being deprecated as of Kafka 3.0, and won't be dropped until Kafka 4.0. Now, if you want to use new shiny features (like cooperative sticky assignors to minimise consumer group stop the world rebalance pauses), then yes, you might need to upgrade clients. But otherwise, you can still happily use 0.8 clients with your upgraded brokers. https://cwiki.apache.org/confluence/display/KAFKA/KIP-724%3A+Drop+support+for+message+formats+v0+and+v1 https://cwiki.apache.org/confluence/display/KAFKA/KIP-724%3A...
- EdwardDiego 5y agoWhat was the bug, out of curiosity?
- at0mic22 5y agoShould you not consider every article as a divine insight. Sixfold is nowhere close from being a classy tech comp.
- anotherhue 5y agoI ran a few dozen kafka clusters at MegaCorp in a previous life. My answer to anyone who asks for kafka: Show me that you can't do what you need with a beefy Postgres.
- bsaul 5y agousing a sql db for push/pop semantic feels like using a hammer to squash a bug.. How would you model queues & partitions with ordering guarantees with pg ?
- anotherhue 5y agosql has many conveniences for doing so, it wouldn't be much work. > using a hammer to squash a bug.. Agreed - but Kafka is a much much bigger hammer. SES/Az Queues are also good choices.
- throwaway81523 5y agoWith transactions, and stored procedures if that helps ;). Redis also seems well suited to the use cases I've seen for Kafka. Kafka must have capabilities beyond those use cases, and I've sometimes wondered what they are.
- deleted 5y ago[deleted]
- afandian 5y agoI had a great time with Kafka for prototyping. Being able to push data from a number of places, have mulitple consumers able to connect, go back and forth though time, add and remove independent consumer groups. Ran in pre-production very reliably too, for years. But for a production-grade version of the system I'm going with SQL and, where needed, IaC-defined SQS.
- aqme28 5y ago> Show me that you can't do what you need with a beefy Postgres. I've found this question very useful when pitched any esoteric database.
- bsaul 5y agoVery interested to hear how people here overcome the limits of kafka for ordered events delivery in real world, and what those were.
- luxurytent 5y agoI feel as if you're using Kafka and expect guaranteed ordering, then you're using the wrong tool. At best you have guaranteed ordering per partition but then you've tied your ordering/keying strategy to the amount of partitions you've enabled ... which may not ideal. But, that's speaking from my light experience with it. I'm also curious if there's a better way :-)
- orobinson 5y agoAt lower data volumes (<10,000 events per minute) it’s perfectly feasible to just use single partition topics and then ordered event delivery is no problem at all. If a consuming service has processing times that means horizontal scaling is necessary then the topic can be repartitioned into a new topic with multiple partitions and the processing application can handle sorting the outputted data to some SLA.
- BFLpL0QNek 5y agoIt depends on what the events are, how they are structured. You get guaranteed ordering at the partition level. Items are partitioned by key so you also get guaranteed ordering for a key. If you have guaranteed ordering for a key you can’t get total ordering across all keys but you can get eventual consistency across the keys. Ultimately if you want ordering you have to design around being eventually consistent. I don’t read a lot of papers but Leslie Lamports Time, Clocks, and the Ordering of Events in a Distributed System gave me a lot of insight in to the constraints. https://lamport.azurewebsites.net/pubs/time-clocks.pdf https://lamport.azurewebsites.net/pubs/time-clocks.pdf
- sumtechguy 5y agoFor kafka the default is round robin in each partition. A hash key can let you direct the work to particular partitions. Each partition is guaranteed ordering. Also only one consumer in a consumer group can remove an item from a partition at a time. No two consumers in a consumer group will get the same message.
- fafle 5y agoThe issue of running a transaction that spans multiple heterogeneous systems is usually solved with a 2 phase commit. The "jobs" abstraction from the article looks similar to the "coordinator" in 2PC. The article does not talk about how they achieve fault tolerance in case the "job" crashes inbetween the two transactions. Postgres supports the XA standard, which might help with this. Kafka does not support it.
- eternalban 5y agoIIRC ~2 decades ago we were dequeueing from JMS, updating RDBMS, and then enqueuing all under the cover of JTA (Java Transaction API) for atomic ops. https://docs.oracle.com/en/middleware/fusion-middleware/12.2.1.4/ashia/jms-and-jta-high-availability.html#GUID-C4750828-F83A-442A-8BFE-EE2F2D3945EE https://docs.oracle.com/en/middleware/fusion-middleware/12.2... Using a very broad definition of ‘noSQL’ approach that would include solutions like Kafka, the issue becomes clear: A 2PC or ‘distributed transaction manager’ approach ala JTA comes with a performance/scalability cost — arguably a non-issue for most companies who don’t operate at LinkedIn scale (where Kafka was created).
- jgraettinger1 5y agoI can't speak to their solution, but when solving an equivalent problem within Gazette, where you desire a distributed transaction that includes both a) published downstream messages, and b) state mutations in a DB, the solution is to 1) write downstream messages marked as pending a future ACK, and 2) encode the ACK you _intend_ to write into the checkpoint itself. Commit the checkpoint alongside state mutations in a single store transaction. Only then do you publish ACKs to all of the downstream streams. Of course, you can fail immediately after commit but before you get around to publishing all of those ACKS. So, on recovery, the first thing a task assignment does is publish (or re-publish) the ACKs encoded in the recovered checkpoint. This will either 1) provide a first notification that a commit occurred, or 2) be an effective no-op because the ACK was already observed, or 3) roll-back pending messages of a partial, failed transaction. More details: https://gazette.readthedocs.io/en/latest/architecture-exactly-once.html https://gazette.readthedocs.io/en/latest/architecture-exactl...
- jgraettinger1 5y agoIf you're in the Go ecosystem, Gazette [0] offers transactional integrations [1] with remote DB's for stateful processing pipelines, as well as local stores for embedded in-process state management. It also natively stores data as files in cloud storage. Brokers are ephemeral, you don't need to migrate data between them, and you're not constrained by their disk size. Gazette defaults to exactly-once semantics, and has stronger replication guarantees (your R factor is your R factor, period -- no "in sync replicas"). Estuary Flow [2] is building on Gazette as an implementation detail to offer end-to-end integrations with external SaaS & DB's for building real-time dataflows, as a managed service. [0]: https://github.com/gazette/core https://github.com/gazette/core [1]: https://gazette.readthedocs.io/en/latest/consumers-concepts.html#stores https://gazette.readthedocs.io/en/latest/consumers-concepts.... [2]: https://github.com/estuary/flow https://github.com/estuary/flow
- ivanr 5y agoSmall suggestion: If Gazette is ready for a wider adoption, it may be useful to bump it up to 1.0 as a signal of confidence.
- LgWoodenBadger 5y agoI didn't quite follow their explanation for why producing to Kafka first didn't/wouldn't work for them (db state potentially being out of sync requiring continuous messaging until fixed).
- anentropic 5y agoit's a chicken and egg problem you can either send a kafka message but potentially not commit the db transaction (i.e. an event is published for which the action did not actually occur) or commit the db transaction and potentially not send the kafka message it sounds like they implemented something like the Transactional Outbox pattern https://microservices.io/patterns/data/transactional-outbox.html https://microservices.io/patterns/data/transactional-outbox.... i.e. you use the db transaction to also commit a record of your intent to send a kafka message - you can then move the actual event sending to a separate process and implement at-least-once semantics This is the job queueing system they described in the article
- LgWoodenBadger 5y agoTheir solution seems like a "produce to Kafka first" but with extra steps. Regarding: When we produce first and the database update fails (because of incorrect state) it means in the worst case we enter a loop of continuously sending out duplicate messages until the issue is resolved I don't understand where either 1) the incorrect state or 2) the need to continuously send duplicate messages come from. Regarding: The Job might still fail during execution, in which case it’s retried with exponential backoff, but at least no updates are lost. While the issue persists, further state change messages will be queued up also as Jobs (with same group value). Once the (transient) issue resolves, and we can again produce messages to Kafka, the updates would go out in logical order for the rest of the system and eventually everyone would be in sync. This is the part that is equivalent to Kafka-first, except with all the extra steps of a job scheduling, grouping, tracking, and execution framework on top of it.
- anentropic 5y agoWhen we produce first and the database update fails (because of incorrect state) it means in the worst case we enter a loop of continuously sending out duplicate messages until the issue is resolved the article does not explain things very clearly, but I think this is describing the problem rather than their solution Our high level idea was: - Insert “work” into a table that acts like a queue - “Executor” takes “work” from DB and runs it ... A Job is an abstraction for a scheduled DB backed async activity ... How did we solve the #2 state problem? By recording Jobs in the service database we can do the state update within the same transaction as inserting a new Job. Combining this with a Job that produces the actual Kafka message, allows us to make the whole operation transactional. If either of the parts fails, updating the data or scheduling the job, both get rolled back and neither happens. I think this is describing basically a Transactional Outbox i.e. "jobs" are recorded in the postgres db as part of the same db transaction as the business logic actions the difference from Kafka-first is that if the app decides to rollback the business logic then the message hasn't already sent
- mfateev 5y agotemporal.io provides much higher level abstraction for building asynchronous microservices. It allows one to model async invocations as synchronous blocking calls of any duration (months for example). And the state updates and queueing are transactional out of the box. Here is an example using Typescript SDK: async function main(userId, intervals){ // Send reminder emails, e.g. after 1, 7, and 30 days for (const interval of intervals) { await sleep(interval * DAYS); // can take hours if the downstream service is down await activities.sendEmail(interval, userId); } // Easily cancelled when user unsubscribes } Disclaimer: I'm one of the creators of the project.
- sam0x17 5y agoCouldn't the "state" issue be solved simply by enclosing the database save and kafka message send in the same database transaction block and only doing the kafka send if it reaches that part of the code?
- staticassertion 5y ago(1) seems best solved by having a 'on heavy task, publish to a secondary topic'. This is good if you have flaky messages that need to be retried in the background, without blocking all of your 'good' messages. (2) this problem should be avoided in general by just having idempotent services. Just as a hard restriction, forever, build services to be idempotent. It should be the exception to have a non-idempotent service, and it should be carefully understood. That said, if you have (1) as a consistent issue, like if every message is flaky, kafka isn't the right solution. Postgres-based queues are perfect for this because you can examine the table as a whole, making more informed decisions about what you want to process (or not process).
- rajin444 5y agoThis seems like a much simpler solution. The "heavy task queue" would process tasks that could simply be retried until they're done. Maybe I'm misunderstanding the article, but having "Job tasks" both insert another Job to run as well as updating DB state, and then having the executor pick up the previously inserted Job (whos only purpose is to send a kafka message) seems overly complex. I'm having trouble seeing why this is needed.
- derekperkins 5y agoSolving both these problems is the best hidden feature of Vitess - Messaging. You can ack a message, do data work, and add to a destination queue all in a single transaction. You can do selective requeuing, and infinite parallelization, since it isn't sequential processing. It's easy to introspect, debug, and have metrics for since it's just MySQL. Vitess lets you shard horizontally, so it can handle any QPS, and has at YouTube. It also supports native message priority. All of that, plus your infrastructure is simplified because you don't have to maintain a separate data store from your main RDBMS. Highly recommended. https://vitess.io/docs/reference/features/messaging/ https://vitess.io/docs/reference/features/messaging/