Hi everyone, I am trying to evaluate Kafka + Pinot...
# general
j
Hi everyone, I am trying to evaluate Kafka + Pinot as a replacement for OpenSearch on a 15.5bn document dataset. There are a number of aggregations we run on the events that group by a user id and then run some statistics (min, max for a number of timestamps, counts of events that match additional specific criteria, etc) as well as some text based analytics. I think solving the statistics part of the aggregations will be straightforward (if very different from what we are doing in opensearch). The core of my question is whether there is any capability for running an aggregation similar to Significant Terms in the Kafka / ksqldb / Pinot suite of tools. What we are trying to produce is a list of the top X results from a given field (email addresses, url domain names, etc) that are most unique in this user's subset of data, compared to all data.
k
Significant terms feature looks pretty cool but seems like it can result in lot of false positives..interesting to hear your take on this feature. If there is a good usecase, adding this feature is not hard.. internally, it’s just doing two or three queries and computing the ratio. And you can do it yourself on the client side as a work around.
j
From my perspective it is a very useful feature to have. There is definitely some nuance to using and understanding it correctly, and tuning the aggregation to provide useful results. In our opensearch workflows we do some pre-processing of the fields that we run significant terms on through the analyzer facilities.
I have read the details on the scoring algorithms for Significant Terms, and while the math isn't necessarily hard, its rather tedious to have to setup multiple queries and do the math as a user when you want to run the aggregation against multiple fields and when you want to group events differently and then run the aggregation. I believe the values needed would already be directly available from the text indexes, whatever is required to calculate cardinality of the data in a field. At an abstract level I guess it is a different flavor of a cardinality type aggregation
k
Agreed on the additional overhead for the user.. any idea on how a sql syntax would look like for this feature?
j
Absolutely no idea. I've not written SQL in almost 10 years. I'm used to building queries with JSON / Java APIs.
If there are other aggregation functions that return a list rather than a scalar, that could be a model to follow.
or maybe this:
Copy code
SELECT user_id, url_host, SIG_TERMS_SCORE(url_host) 
FROM  events
GROUP BY user_id, url_host 
HAVING SIG_TERMS_SCORE(url_host) > 90 
ORDER BY SIG_TERMS_SCORE(url_host) DESC 
LIMIT 50
since what we really need is the score
I think all we might need to do this is the score. I was thinking about it too much like OpenSearch.
k
group by user_id?
j
let me clean that up a bit...
in an easy to explain version of our use case, we have data streaming in that has a user_id field, url fields, and about 20 - 30 other fields. We want to group our data by user first, and then see the terms from one of our text fields that are the most unique (have the highest Significant Terms score) for each user. I think in proper SQL syntax we would need to group by that text field as well. You will have to excuse any bad SQL in my example. Hopefully the description above makes the idea more clear.