4 ms·
AIUI you can only write to the database indirectly through Kafka messages which the KSQL then periodically consumes and materializes into tables with syntax alm
by Erwin 7y ago
AIUI you can only write to the database indirectly through Kafka messages which the KSQL then periodically consumes and materializes into tables with syntax almost like SQL.
So if you emit an event (like "New User Created") and have a ksql table that summarizes the user count, then Kafka having accepted the event does not mean the ksql tables have also been updated.
Contrasted with a traditional database where a COMMIT returning one one session means the second can immediately read it.
- missosoup 7y agoThat's standard behaviour in every analytics database that talks to kafka that I can think of. Kafka doesn't really encourage ACK/NACK type processing patterns. Accepting an event usually means the consumer has successfully read it and staged it for whatever is meant to happen next, not that the operation is completed. Now if it were possible for it to accept an event from Kafka but not guarantee that the event will eventually make it into the materialized view but may be lost, that'd be a problem.
- bladecatcher 7y ago> Now if it were possible for it to accept an event from Kafka but not guarantee that the event will eventually make it into the materialized view but may be lost, that'd be a problem. This is usually not a problem these days as it’s possible to guarantee exactly once ingestion using Kafka offsets
- missosoup 7y agoIIRC clickhouse still doesn't have this guarantee. It guarantees exactly one ingestion but not that the ingested event gets processed all the way through to view. And if the processing fails, that event won't be retried and is now gone.