[1] https://github.com/apache/flink/blob/master/flink-connectors... [2] https://github.com/apache/flink/blob/master/flink-connectors...
21 karma · joined October 9, 2017
[1] https://github.com/apache/flink/blob/master/flink-connectors... [2] https://github.com/apache/flink/blob/master/flink-connectors...
I think the space for streaming processing is still quickly evolving. Many features like stream-table joins, CTEs, streaming joins with late arrival data are unimplemented or do not even have clear semantics yet. It would be great to see a benchmarks like TPC-DS in the domain.
(1) There are significant loads on consultations when users had to implement their own jobs in Java / Scala and run them in production. Sometimes it turned in to co-development as the users lack the expertise of the streaming analytics frameworks.
(2) We consciously encourage our users to write good SQLs via: (a) enforcing schemas on all analytical Kafka topics. (b) setting up a team dedicated to help them using SQL in big data systems (i.e., Hive, Presto, AthenaX, etc.)
For UDF we provide general guidances and ask our users to oncall for the jobs that use UDFs. The support costs are definitely not zero but it is still much better to teach users to write a Samza / Flink / Storm job from scratch.
Given that Kafka neither provides secondary indexes nor organizes data in columnar format, Presto essentially have to somewhat scan through the Kafka topics to execute the queries, resulting a lot of disk I/O.
Our Kafka infrastructure handles more than one trillion messages per day and guarantees second-level latency SLAs. Reading aggressively could easily saturate all the I/O bandwidths of the nodes and leads to outages. We actually had several incidents in the past when we did backfills. So I'm more conservative on this.
What do you think of posting the questions on the mailing list [1]?
We internally have a React-based UI and we are in the process of cleaning it up and opening sourcing it. Please stay tuned!
But yes streaming analytics is a very exciting area. Your feedbacks are appreciated.
Proposals on names are always welcomed.
Personally I'm not a big fan of KSQL, given that:
(1) KSQL is built on top of Kafka Stream. There are use cases that we don't think Kafka Stream is a good fit. Please see the explanations above
(2) Inventing yet another SQL dialect is a bad idea in practice. Not only it incurs additional learning curves, but more importantly you have few hopes winning in development velocity. Calcite has been used by Flink, Storm, Beam, Dremio, etc. The community is simply way bigger even compared to the total number of engineers in Confluent.
Again this is just my personal take it does not reflect the stands of Uber.
Before AthenaX most of our real-time analytic pipelines were on Samza, which to to some extent can be seen as predecessor of Kafka Stream.
The migration from Samza towards Kafka Stream might seem more natural, and we actually took a very close look on Kafka Stream and we have decided to move towards Flink, given that:
(1) Kafka Stream lacks of important features like exactly-once delivery and distributed consistent snapshots. They are essential to support use cases that require high fidelity.
(2) SQL is a must-have feature in order to empower our users, many of which are non-technical, to run large-scale streaming analytics in production. Kafka Streams provide no support for that. Arguably the Kafka Streams provide simple APIs -- Simple APIs themselves are insufficient to bring the analytics applications to production. You can't ignore continuous integrations and deployment, monitoring, etc., especially many of our users come from non-CS backgrounds.
(3) It seems that the Apache Flink community is more open and committed compared to the Kafka community, particularly on the SQL side. We have collaborated with Data Artisans, Alibaba and Huawei. All parties above, as well as us are committed and equipped with adequate resources to bring SQL to respective customers.