Hi, I’m trying to improve the performance of a sel...
# troubleshooting
a
Hi, I’m trying to improve the performance of a select count(distinct col_a) query, it’s taking several minutes at the moment before failing (out of memory, box has 64gb ram). There are about 50 million unique values from about 700 millions rows. The DistinctCountHLL and DistinctCountThetaSketch estimates are fast enough but not accurate enough. What can I do improve the performance of the count(distinct col_a) query?
r
what's the type of the values?
a
string
the query will often have a group by and where clause in it, these tend to be string values as well
r
it's kind of a hard problem then, even if the strings are relatively short, say 20 bytes, 50M of them requires ~1GB RAM to deduplicate
do you really need an exact result?
a
Ideally, yes. The estimates are off by 10's of thousands, often by 100's of thousands, which is too much.
Can I use any of the indexes in Pinot to help? The data doesn't change very often, a few times a month, so willing to accept longer ingest time for better query performance.
r
can you share the table config for the field you are running distinct count over?
I think we can probably improve distinct count to use indexes and dictionaries where available, but I don't think this is possible right now
a
A specific example of the numbers I'm getting at the moment (with a group by on another column) Actual 15476256 distinctCountThetaSketch 15327774 DistinctCountHLL 14866424
Copy code
{
  "tableName": "op_test",
  "tableType": "OFFLINE",
  "segmentsConfig": {
    "segmentPushType": "APPEND",
    "segmentAssignmentStrategy": "BalanceNumSegmentAssignmentStrategy",
    "replication": "1"
  },
  "tenants": {},
  "tableIndexConfig": {
    "enableDefaultStarTree": true,
    "loadMode": "MMAP",
    "invertedIndexColumns": [
      "gender"
    ]
  },
  "metadata": {
    "customConfigs": {}
  }
r
I think you can do better with theta sketch by setting the nominal entries in the query
a
Could explain with an example please?
r
by default nominal entries is 4096, but you can set it like so:
Copy code
DISTINCT_COUNT_THETA_SKETCH(col, 'nominalEntries=8192')
and you can go even higher
so basically you should be able to get more accurate results by increasing it, but it will use more memory, though not as much as an exact distinct count
can you give it a try for a few values (e.g. 8192, 16384) and see if you get better results?
a
'nominalEntries=16384', gives a count of 15426624, takes about 15 seconds.
r
ok, the size of the sketch increases with that parameter so an increase in latency is to be expected
a
actual: 15476256 'nominalEntries=65536', count = 15472132 'nominalEntries=131072', count = 15489003
still takes about 15 seconds
r
how long was it without specifying nominal entries?
a
the same, about 15 seconds without nominal entries
I wasn't expecting 'nominalEntries=131072', to be less accurate than 'nominalEntries=65536'.
r
I think that's larger than the sketch authors really intended, there may be bugs assuming you won't go that large
a
Ah
r
what are the strings being counted? Could they be transformed to numbers on ingestion?
a
I guess they could be. They are IDs...example:
Copy code
9IPX6Q4ZD9TT8QD
K7IL68L885FWT07
HINHAWWUCYUADZB
r
you're getting pretty good relative error from theta sketch, ((15476256 - 15472132) / 15476256) * 100 = 0.03%
at 16384 ((15476256 - 15426624) / 15476256) * 100 = 0.3%
a
ideally the results need to be accurate to less than 5
r
most sketches don't have absolute error guarantees though
is the column with the IDs raw or does it have a dictionary?
a
What do you mean?
r
can you share the schema definition for the ID column?
I think you should try sorting on the ID column and see where that gets you
a
Copy code
....{
      "name": "heid",
      "dataType": "STRING"
    },...
r
so
Copy code
{
  "tableName": "op_test",
  "tableType": "OFFLINE",
  "segmentsConfig": {
    "segmentPushType": "APPEND",
    "segmentAssignmentStrategy": "BalanceNumSegmentAssignmentStrategy",
    "replication": "1"
  },
  "tenants": {},
  "tableIndexConfig": {
    "enableDefaultStarTree": true,
    "loadMode": "MMAP",
    "invertedIndexColumns": [
      "gender"
    ],
    "sortedColumn": ["heid"]
  },
  "metadata": {
    "customConfigs": {}
  }
a
Once I've made the config change, I'll have to ingest the data again?
r
yes, and I'm not 100% sure it's going to help, but I think it might because what the distinct count aggregation function does is build a compressed bitmap from the internal integer representation of
heid
however, if you ever want to filter on
heid
sorting on it will really help, it's generally good to sort on high cardinality attributes
a
Thanks Richard, I'll try it out. Will take a few hours to ingest the data.
r
if you can, take a profile of the query with
jcmd <server pid> JFR.start duration=60s filename=distinctcount.jfr
and send it to me and I'll take a look to be sure what I think is happening is happening, but do that before reingesting
I also don't want you to spend hours doing that if I'm wrong, and having a profile would help
a
That's fine, I'm doing a PoC on Pinot and the ingest can run while I'm doing other things.
On a side note, when I set
segmentCreationJobParallelism: 14
from 1, some of the segments/data is missing when the ingest is completed. I couldn't find anything in github issues about it...nor any errors in the logs. Any ideas what could be going wrong?
r
Can we create another thread for that? I'm not too hot on ingestion topics
✅ 1
by the way do you have the query and the response metadata? Replace the aggregation function with count so it comes back quicker, maybe adding some indexes will help.
k
If you really need distinctCount accurate and also speed.. you can partition data by these ids and use partitioneddiatinctcount
Partitioned DistinctCount
a
@Richard Startin, could you explain what you meant by 'Replace the aggregation function with count so it comes back quicker, maybe adding some indexes will help.'
@Kishore G, how do I setup Partitioned DistinctCount, could you point me to the docs for it? Thanks
r
rather than putting a distinctcount in the query, put a count, looking at the response metadata would be helpful to figure out if you need some indexes
the biggest problem right now will be that exact distinctcount is such an expensive thing to compute though
a
The count comes back fairly quick from what I can recall. I get back to you with metadata once the data has been ingested.
k
@Jackie ^^
j
What's the latency of
distinctCountHll
?
You can increase the
log2m
for
distinctCountHll
to get better accuracy, e.g.
distinctCountHll(col_a, 15)
You can also try
distinctCountBitmap(col_a)
which counts the accurate unique hash values with java
String.hashCode()
a
Copy code
redshift
count: 15476256
time: 525.429 seconds

select distinctCountHll(heid) 
timeUsedMs: 18107
count: 14866424

select distinctCountHll(heid, 15) 
timeUsedMs: 17140
count: 15375247

select distinctCountHll(heid, 30) 
OutOfMemoryError

select distinctCountBitmap(heid) 
timeUsedMs: 43089
count: 15448390
the table config for this table is slightly different than previous:
Copy code
"tableIndexConfig": {
    "loadMode": "MMAP",
    "invertedIndexColumns": [
      "gender",
      "heid"
    ]
  },
the response metadata your referring to above, is that just the response json when running the query through the pinot query console?
k
Jackie, @Ali is looking for partitioned distinct count function. I remember building this for folks from Target, do we have any docs for that?
a
@Richard Startin, regarding the jcmd, I was using the slack web client earlier which didn’t render some of your message. also, the data I’m querying is sensitive, not sure if the dump will contain any of the data.
j
@Ali here you can find the documentation of the partitioned distinct count function: https://docs.pinot.apache.org/users/user-guide-query/supported-aggregations
It is very fast and can get accurate result, but requires the column to be partitioned per segment (i.e. no common value across segments).
I feel
select distinctCountHll(heid, 15)
is giving quite close result with low latency
r
@Ali JFR was designed to send diagnostics back to Oracle and avoids any kind of PII for the reasons you're concerned about
the method profiling tab would show you where the time is being spent
a
Hi, I’ve reingested the data with the following config:
Copy code
"tableIndexConfig": {
      "segmentPartitionConfig": {
        "columnPartitionMap": {
          "hesid": {
            "functionName": "Murmur",
            "numPartitions": 32
          }
        }
      },
      "sortedColumn": [
        "hesid"
      ]
When I do a distinct count by hesid, I’m getting an OutOfMemoryError, the timeout is 120seconds.
r
so it OOMS before 120s?
@Mayank
m
What’s the exact query and total doc count?
a
Yeah, the OOMS is before the 120 seconds. TotalDocs is 635904509, the query is
Copy code
select count(distinct hesid), financial_year 
from op_test 
group by financial_year
limit 10
m
This is an expensive query and you probably have less resources to compute. Do you need exact results or are you ok with fast approximate? If yes, you can try distinctcountHll.
r
we discussed that last week, the relative error is less than 0.5% and the queries are fast, but Ali needs absolute error of +/- 5
m
Then the only other optimization I can think of is partitioned way of count distinct: https://github.com/apache/pinot/pull/5786
a
When I change the query to select to
Copy code
select SEGMENTPARTITIONEDDISTINCTCOUNT(hesid), financial_year 
from op_test 
group by financial_year
limit 10
the query comes back within a few seconds but the numbers are too small e.g. 36159185 instead of the actual 49233971 for a particular financial_year
r
ok so looking at the description, this function won't work for your use case because it assumes partitioning per segment
so this would require a global sort on
hesid
for it to work properly... exact distinct counts on large data sets are hard...
a
Copy code
"tableIndexConfig": {
      "segmentPartitionConfig": {
        "columnPartitionMap": {
          "hesid": {
            "functionName": "Murmur",
            "numPartitions": 32
          }
        }
      },
      "sortedColumn": [
        "hesid"
      ]
This is the config I have setup, the latter part of the config is a global sort on the hesid isn’t it? and how do I setup partitioning per segment.
As I explained earlier, my data doesn’t change regularly, so anything I can do at ingest to speed things up is worth considering
r
the problem with segment level partitioning on
hesid
is that if you ingest some more data and some `hesid`'s you've seen before show up, they would need to go in to already built segments, otherwise the partitioning would be violated by the new data
so you would need to rebuild those segments at that time
k
@Richard Startin thats not true, you can have multiple segments with all keys belonging to same partitionId
what this feature requires is that it the right assignment strategy so that all segments belonging to give partition are assigned to same instance
we already have such as assignment strategy, it needs to be configured in table config.. @Jackie ^^
a
Could someone suggest how this config should be adjusted to do what you’ve suggested then:
Copy code
"tableIndexConfig": {
      "segmentPartitionConfig": {
        "columnPartitionMap": {
          "hesid": {
            "functionName": "Murmur",
            "numPartitions": 32
          }
        }
      },
      "sortedColumn": [
        "hesid"
      ]
j
@Richard Startin is correct, this function requires segment level partitioning, instance level is not enough. But if the segments are not partitioned, we should get larger result instead of smaller, so there must be something else wrong
If 100% accurate result is not required, distinctcouthll should give the best performance, and you can tune the accuracy using the optional log2m as the second argument
r
it is required
as mentioned several times in the thread, Ali would accept +/- 5 absolute error, even 0.03% relative error was too much with theta sketch (we saw similar relative error with hll)
I have profiles of this and it's very clear what's going on execution wise, the question is what does @Ali need to do to achieve segment level partitioning on
hesid
, and if he ingests new data with some of the same
hesid
what does he need to do to ensure the partitioning isn't later violated, assuming he queries over large time ranges
m
Since the thread has gotten too long, let me summarize:
Copy code
1. The query is performing accurate count distinct on 635M records, which is quite an expensive operation. (Not sure if the thread mentioned what's the resource provided to servers).
2. PartitionedDistinctCount currently only works when partition key is limited to a segment. So either it has to be a REFRESH use case, or a case where keys are scoped within a day (unlikely?).
3. For such cases we recommend distinctCountHLL which provides high accuracy approximation at low latency.
a
Thanks for summarising, the data changes about 3/4 times a month and I think we’ll have to accept doing a full ingest if that is what is required. The server is 16 cores, 64MG ram, on AWS.
the
distinctCountHLL
is fast but not accurate enough, it’s off by 10's of thousands and I need something either exact or to the nearest 5.
m
What’s the total data size?
a
about 700gb
m
Ok, if you are ok with data refresh then you could use partitioned approach
a
Great, so I have the following config at the moment:
Copy code
"tableIndexConfig": {
      "segmentPartitionConfig": {
        "columnPartitionMap": {
          "hesid": {
            "functionName": "Murmur",
            "numPartitions": 32
          }
        }
      },
      "sortedColumn": [
        "hesid"
      ]
How do I change it to use the ‘partitioned’ approach?
m
If pushing offline, you have to partition the data before indexing and pushing to Pinot. And then enable partition based segment assignment https://docs.pinot.apache.org/operators/operating-pinot/instance-assignment#partitioned-replica-group-instance-assignment
a
Thanks, I’ll have a look.