Hello! I’m running PoC with Pinot for quite heavy ...
# troubleshooting
p
Hello! I’m running PoC with Pinot for quite heavy data. I’ve got table with ~ 20 billions rows, 5 predicates (pred1 cardinality ~ 5 millions uniqs, pred2 and pred3 the same - tight correlation, last pred5 has very low cardinality, tens of values). I need to achieve the best possible speed for lookups by these predicates for whole range (20-50 billions/rows). Currently my table (realtime) is creating for these predicates bloom & inverted indices. Second problem is ingestion rate - apparently there is no problem to get ~ 160k/s documents which is insane in contrast to resources needed, but at the same time the query performance is very bad - 6 servers are pretty busy with ingesting and GC thus query is pretty bad, 20-50 seconds. My current setup is 6 servers, 1 controller. Split commit enabled to s3. Because there will be low QPS, I need to achieve low memory allocation for indices/segments. Do I need to consider some kind of bucketing/hidden partitioning for predicate values or is Pinot able to handle these data in SLA ~ 1000-3000 ms only with proper indexing? I can imagine some sort of work delegation for servers, e.g. consuming/segment creating ~ 3-4 servers and for querying allocate 6 servers. PS: I’ve got replication 1 for space saving as final total will be ~ 20 TB, segment size is currently 460MB (but in table is set to 1GB). Ingesting from 36 kafka partitions. Servers have spin drives. PS2: not sure about effect of bloom because is still in rebuild process… Any improvements, thoughts or tricks are welcomed! 🙂
r
first question: which version of pinot are you using?
p
0.9.3
I’m sure if my expectations are even realistic…
r
to understand the slow queries, it would be good to get a profile during query execution - you can collect this by running
jcmd <realtime server pid> JFR.start duration=60s settings=profile filename=slow_realtime_query.jfr
next, knowing your ingestion config (e.g. what kind of transforms are you running?) the indexing config (which indexes do you have?) and an obfuscated slow query would be helpful
also, knowing something about the data would be good to help understand how selective the indexes you have set up are - e.g.
select indexed_column,count(*) as cnt from table group by  indexed_column order by cnt desc
the response metadata on the queries (JSON view in the query console) will tell you how many docs were scanned in different stages of the query, please paste that here for a slow query
p
I’ll capture flight profile later, but regarding to other questions: 1. ingestion - only flattening of json file (e.g. 8 filelds, via JSONPATH…, one field kept as “rich” JSON object in STRING see tests_table config attached 2. Cardinality is described in json - 4x predicate with ~ 6 millions unique values, 1 predicate with cca 60 unique values 3. Other responses/traces I’ll post when cluster will be ready. I’m suspecting that kubernetes cluster has additional issues.. Running on pretty strong bare metal, but old kernel, sometimes I see broken pipe and network connectivity issues between Pinot components
Defined resources: Controller 1x: resources: requests: memory: “8Gi” cpu: “8000m” limits: memory: “16Gi” cpu: “16000m” jvmOpts: “-Xms1G -Xmx6G -XX:+UseG1GC -XX:MaxGCPauseMillis=200 -Xlog:gc*:file=/opt/pinot/gc-pinot-controller.log -javaagent/opt/pinot/etc/jmx prometheus javaagent/jmx prometheus javaagent 0.12.0.jar=8008/opt/pinot/etc/jmx_prometheus_javaagent/configs/pinot.yml” Server 6x: resources: requests: memory: “32Gi” cpu: “12000m” limits: memory: “48Gi” cpu: “32000m” jvmOpts: “-Xms1G -Xmx16G -XX:+UseG1GC -XX:MaxGCPauseMillis=200 -Xlog:gc*:
r
if you're doing jsonpath I recommend trying the latest docker images on 0.10.0-SNAPSHOT - I fixed a lot of performance problems in the jsonpath lib
👍 1
it's not perfect, but it's better
looking at your config, you don't have a sorted column, and that would help, choose one of
pred1-4*6m
for sorted
👍 1
whichever you query most by
are any of the pred columns numeric?
p
Unfortunately all strings.. but I think that the mostly used predicated could be converted to long
r
using an onheap dictionary for pred5_cardinality_60 (and only pred5_cardinality_60) might be helpful to save on allocations
the strings get dictionarized anyway, and will use log_2(6m) = 23 bits per value when converted to offline, but converting to long beforehand will help for real time, if you can
I'll wait for the query, response metadata, and the profile (I'm fairly sure JSONPath will figure in the profile heavily)
p
Ok, and what about inverted index? Is it even helpful here in combination with bloom? Or bloom would be enough for lookup?
r
it would be good to know what difference you get from the nightly docker images before reconfiguring anything
p
Just upgrading 🙂 should be running already
r
hard to say, need to see the response metadata and a profile really
inverted indexes depend on the distribution, if it's uniform, 6m bitmaps should be highly selective, but if 50% of the data has
pred1_cardinality_6m == x
for some x, then they're not so good
p
at least per segment is uniform, with high selectivity
r
ok good, let's just wait for the data to see what's going wrong - your expectations are realistic - but some features are more mature than others
p
Is it OK to capture flight profile from 1 server (from 6)?
r
yes, it should be evenly split
p
Flight profile - search by mostly used predicate, response in 50 sec. Currently table is “only” 7 billions.. but without already mentioned improvements such as sort etc..
r
thanks
do you also have the response metadata
p
yes, it’s for different value but now with time 126 sec “exceptions”: [], “numServersQueried”: 6, “numServersResponded”: 6, “numSegmentsQueried”: 3094, “numSegmentsProcessed”: 1872, “numSegmentsMatched”: 109, “numConsumingSegmentsQueried”: 36, “numDocsScanned”: 162, “numEntriesScannedInFilter”: 0, “numEntriesScannedPostFilter”: 7614, “numGroupsLimitReached”: false, “totalDocs”: 7028875200, “timeUsedMs”: 126651, “offlineThreadCpuTimeNs”: 0, “realtimeThreadCpuTimeNs”: 0, “offlineSystemActivitiesCpuTimeNs”: 0, “realtimeSystemActivitiesCpuTimeNs”: 0, “offlineResponseSerializationCpuTimeNs”: 0, “realtimeResponseSerializationCpuTimeNs”: 0, “offlineTotalCpuTimeNs”: 0, “realtimeTotalCpuTimeNs”: 0, “segmentStatistics”: [], “traceInfo”: {}, “minConsumingFreshnessTimeMs”: 1645106201566, “numRowsResultSet”: 10
r
and the query (redact the params)?
p
select * from epc_pgw where a_party_imsi = ‘xxxx’ limit 10
✅ 1
r
virtually all the time is spent on processing JSON
p
Ok it means that ingestion is fully saturating all threads and there is no space to query, right?
r
GC looks fine, only two STW events but thy were under 10ms
👍 1
kafka client is allocating a lot, as usual
at the JVM level, this doesn't look unhealthy
the allocation profile is basically the same as the cpu profile, just with more kafka client showing up
I can't see actually any evidence of queries running from the profile
what happens if you run
select count(*) from epc_pgw where a_party_imsi = 'xxxx' limit 10
- similarly slow?
k
@Pavel Stejskal Query is not doing much, it’s mostly bcos of resource contention between ingestion threads and query threads.. Pinot has the ability to isolate consuming servers and query servers..
You can create x real-time consuming servers and may be y for purely serving and setup the real-time to offline segment task. Cc @Neha Pawar
👍 1
That will allow you to scale the two independently
p
For “select count(*)..” 73 seconds - with stopped ingestion. jfr attached “exceptions”: [], “numServersQueried”: 6, “numServersResponded”: 6, “numSegmentsQueried”: 3099, “numSegmentsProcessed”: 1858, “numSegmentsMatched”: 1765, “numConsumingSegmentsQueried”: 36, “numDocsScanned”: 2742, “numEntriesScannedInFilter”: 0, “numEntriesScannedPostFilter”: 0, “numGroupsLimitReached”: false, “totalDocs”: 7031042881, “timeUsedMs”: 73528, “offlineThreadCpuTimeNs”: 0, “realtimeThreadCpuTimeNs”: 0, “offlineSystemActivitiesCpuTimeNs”: 0, “realtimeSystemActivitiesCpuTimeNs”: 0, “offlineResponseSerializationCpuTimeNs”: 0, “realtimeResponseSerializationCpuTimeNs”: 0, “offlineTotalCpuTimeNs”: 0, “realtimeTotalCpuTimeNs”: 0, “segmentStatistics”: [], “traceInfo”: {}, “minConsumingFreshnessTimeMs”: 1645106324027, “numRowsResultSet”: 1
r
it looks like the servers are just saturated by ingestion
there might be something wrong with that second recording, the recorder itself was blocked on a monitor for 11s, don't worry about that though, I would focus on splitting ingestion from query
what kind of latency do you need?
p
~ 3000 ms. I’ll rebuild segments according your suggestions
r
oh, I don't mean query latency, but end to end latency from event creation to showing up in queries, how fresh does the data need to be?
p
near real-time if possible. I’d say 30 seconds
k
You don’t have to rebuild segments
Add more servers and tag them separately and rebalance
p
It will reorder (sort by co1) automatically?
I changed table config and servers are doing something..
k
I don’t think you need it right now.. let’s validate the separation of consuming from non consuming nodes..
p
I’ve got additional 6 servers, how to tag them for “query” only?
Rebalance exlc. consuming?
See Moving completed segments to different hosts
basically tag servers as _REALTIME or _OFFLINE and then in the table config add the overrides; note that the process to move segments is async so they may not move immediately
👍 2
k
thanks Peter 🙂 I was searching for that and could not find it..
p
One question.. segment itself does contain any index? I see during rebalancing a lot of messages that new server is creating bloom index for just downloaded segment PS: could be only for old segments because I changed indexes later ..
k
Segments in deep store does not have any indexes .. they get created when they are downloaded for the first time
👍 1
You can decide to create them at the time of creation if you want to but it’s not needed bcos you can add delete indexes dynamically
p
And one question to controller.. I think there could be potential bottleneck as well, because responses are pretty slow. I don’t have a dedicated disk for controller and I’m not sure if there is some high I/O traffic. E.g. loading of tables (meta) takes around 15 seconds…
r
can you take a profile of the controller? It can indeed be a bottleneck for segment completion
p
Sure. I'll post details tomorrow when rebalancing will be done
👍 1
Hi, I redeployed whole Pinot cluster because controller was more and more unresponsive - sometime there was null pointers exceptions and errors with combinePlan or mergePlan (something with plan). Currently with fresh deployment I’ve got so far 70 millions rows and latency ~ 10 ms. Actually the same latency as with only few thousands of rows - potentially some “natural” latency because of my infrastructure… I’ll keep you informed when I reach to ~ 20 billions. Anyway, what is the bottleneck for IOps and metadata? Zookeeper must be damn fast (IO), or is there some kind of caching in Pinot components? My controller was yesterday super unresponsive. Is there some recommendation for e.g. balancing controllers? Thanks PS: yesterday, after rebalancing and when everything was idle, query was taking around 50-100 seconds… I’m pretty sure that it was some broken component.. e.g. large metadata on slow disk or similar design issue. Right now I got for zookeeper and controller hostPath PV. PS2: I reduced ingestion rate to ~ 30k/sec (36 kafka partitions to 6) and running on 6 servers, without tag overriding for now…
r
hi - is it still performing well?
getting that profile of the controller if it starts to slow down would be helpful
p
Yes, 220 millions rows and by 1 predicate results ~ 14 ms, always under 20ms. Almost the same as for empty table. There was something weird before. Anyway, helm chart with tag latest could possibly messed my deployment, because when I did a server exapansion, latest was possibly different version and some hosts were downloading images.. I’m not sure if all my images was the same 😞 should be explicitly stated in helm chart
r
hmm, it would be good to get to the bottom of whatever happened, keep us posted if you see the same thing again
👍 1
p
Hello, I’m experiencing the same issue with latency again. Around 1 billion rows everything seems to be fine, p95 ~ 60ms. Close to 2 billions I’m getting more and more timeouts (10sec). jfr file attached witch multiple queries… Running with: 6 servers, each has own spinning 1 drive 8TB, each 48GB memory. The latency is stable till Memory Mapped Usage is aroud 45 GB per server, then is getting worse. Disk has decent performance but is it ordinary spinning drive (120 mb/s, read 400 iops, 120 write iops). Table is currently ~ 1.1 TB (2232 segments).
r
could be page faults
do you have perf or async-profiler on the servers? We could prove it's page faults that way.
if it is page faults causing the slow down, then we could look into enabling hugetlbfs on your servers, which MappedByteBuffers will use
p
No, I don’t have an access to hosts directly, just to kubernetes
r
ok so there's no way you can do kernel level configuration?
you can install linux-perf and async-profiler into the docker image though
p
Yes I can, but I have to describe what to do to my colleague which has access
r
ok, that sounds painful, let's figure out if it would be beneficial first
p
We’re running on latest RHEL but it has kernel 4.18 or something…
r
we run this to get perf
Copy code
apt-get update && \
     apt-get install -y --no-install-recommends linux-base linux-perf openjdk-11-dbg
p
Anyway, is there some limit for MMAP regarding to size of data in total? Per server
r
we install async-profiler like this
Copy code
ARG TARGETPLATFORM
RUN case $TARGETPLATFORM in \
       linux/amd64) arch=x64; ;; \
       linux/arm64) arch=arm64; ;; \
       *) echo "platform=$TARGETPLATFORM un-supported, exit ..."; exit 1; ;; \
     esac \
     && mkdir -p /usr/local/lib/async-profiler \
     && curl -L <https://github.com/jvm-profiling-tools/async-profiler/releases/download/v2.5.1/async-profiler-2.5.1-linux-${arch}.tar.gz> | tar -xz --strip-components 1 -C /usr/local/lib/async-profiler \
     && ln -s /usr/local/lib/async-profiler/profiler.sh /usr/local/bin/async-profiler
generally, yes, there are limits because you only have so much RAM
and by default the page size is 4KB which leads to lots of faults, which hurt more with HDD
what I can see frmo this profile is that the segment is being built, and it's bottlenecked on memory accesses
p
Ok let me understand. If I got server with 48 gigs of Ram, how many segments can this sever hosts in gigs? Rough estimation. To be performing well
r
I'd break it down differently
you need around 16GB for the JVM heap, which will cap direct buffers at 16GB, so thats 32GB from 48GB
that leaves 16GB, the JVM will need another 500MB for native memory, round up to 1GB, which leaves 15GB for the operating system memory mapped buffers and so on
p
Current setup for server: -Xms1G -Xmx16G -XX:+UseG1GC -XX:MaxGCPauseMillis=200, container has limit 48GB
r
what's the segment size in MB?
p
250MB
r
my ex-linkedin colleagues recommended slightly more server memory for a heap of that size based on their production experience - 64GB
p
so xmx=64? Is there some penalty getting over 32 GB (compression…)
r
no I mean xmx=16, server has 64GB
👍 1
you don't want to cross 32Gb because of compressed oops
👍 1
15-20GB feels like a realistic limit here but there's no hard and fast rule, the limit will be determined by when you start paging
would you find publishing that kind of sizing information helpful?
I think you are probably running into page faults which is why you're encountering the wall at 2Bn, proving that would be useful
p
sure 🙂 becauase I’m struggling with basics like this..
r
@Mayank may be able to help later
p
I’ll contact my colleague to install profiler and check hugepage in kernel
r
you need to set
-XX:+UseHugeTLBFS
to make use of it
when Mayank sees this (might be monday) he will be able to help more with sizing, I'd just like to understand what's causing the wall, which will feed into improvements and better resourcing documentation
p
Yes I’d like to understand the limits and bottlenecks as well. But regarding to hosted segments and total size, I still don’t understand what i Pinot capable to handle. If there is limit just disk space, or some special mapping between system memory to segments via MMAP… For example my case, server with xmx=16, 64 gigs ram, how many data it should be capable to handle?
r
There isn't really a simple formula and realtime and offline nodes behave differently. The time taken here is spent in conversion of the realtime segment to an offline segment and having a faster disk or more memory would certainly help with this operation, because it's doing things like translating data into more efficient formats for faster queries as the offline segment will exist for a long time. It's writing to disk, doing compute-heavy operations, buffering. So it's more a question of having the right resources for this particular operation than what the right number of segments is.
one thing that would help here is to review the data types in your schema, since a lot of time is spent sorting strings during the conversion - it's pretty common to start off with strings and fine tune the schema later, but if you have columns which are really ints or longs, that would speed this up.
it's also worth figuring out if all of your fields need dictionaries - do you have any CLOB/BLOB columns? Any JSON? etc.
any numeric columns which you will mostly sum up should be metrics and so on
p
Yep got it. The sorting is for the mostly used predicate, should be converted to long later. Otherwise whole input is json therefore for projection is there few big structures. All columns expect 5 predicates are noDictorinary.. without indexes
None of them are metrics now, its all simple lookups
My question regarding to memory was for offline server, capabilities of hosting segments. If there is some limit or direct correlation between system memory and disk storage. Simply if I can host as many segments as needed, or e.g. 20 TB table needs to be hosted by XY servers, each 64gigs of RAM. Just some simple formula for offline. Realtime is quite complex model for memory calculation… there could be very interesting to mention what impact has single vs multiple kafka consumers per server etc.
m
depending on your latency requirements, you could do 1TB per offline node with 16 cores 64GB, very roughly speaking
p
so basically same ram footprint as data size for offline low latency?
r
my response on size was how much data can fit in RAM on the realtime servers with 48GB
p
Latency p95 ~ 2000-3000 would be fine. Thanks for estimation. I'm ready to design infra exactly for Pinot. We have 396gb ram available per server but I'm not sure about disk drives, how big issue would be spinning drives. If spinning is no go, we can adapt drives as well after PoC.
Same ram as data for low latency - that would be kinda inefficient… but this is my current behaviour and I'm suspecting disk drives.. Or page faults as Richard mentioned
r
I don't think same ram as data for low latency has been suggested, just how much data is served from consuming real time servers.
p
Yes I have. My current state: table with 4,9 billions rows, size 1.85 TB (reported size), 3806 segments. Spread evenly on 6 servers: Server_pinot-server-3.pinot-server-headless.pinot.svc.cluster.local_8098 624 Server_pinot-server-4.pinot-server-headless.pinot.svc.cluster.local_8098 609 Server_pinot-server-5.pinot-server-headless.pinot.svc.cluster.local_8098 665 Server_pinot-server-0.pinot-server-headless.pinot.svc.cluster.local_8098 644 Server_pinot-server-1.pinot-server-headless.pinot.svc.cluster.local_8098 619 Server_pinot-server-2.pinot-server-headless.pinot.svc.cluster.local_8098 645 My query: select * from epc_pgw where a_party_msisdn=‘xxx’ <-- predicate cardinality ~5millions, sorted by this and has bloom + inverted index limit 10 Response: “exceptions”: [], “numServersQueried”: 6, “numServersResponded”: 6, “numSegmentsQueried”: 3806, “numSegmentsProcessed”: 1492, “numSegmentsMatched”: 112, “numConsumingSegmentsQueried”: 6, “numDocsScanned”: 137, “numEntriesScannedInFilter”: 0, “numEntriesScannedPostFilter”: 6439, “numGroupsLimitReached”: false, “totalDocs”: 4902830003, “timeUsedMs”: 19090, “offlineThreadCpuTimeNs”: 0, “realtimeThreadCpuTimeNs”: 0, “offlineSystemActivitiesCpuTimeNs”: 0, “realtimeSystemActivitiesCpuTimeNs”: 0, “offlineResponseSerializationCpuTimeNs”: 0, “realtimeResponseSerializationCpuTimeNs”: 0, “offlineTotalCpuTimeNs”: 0, “realtimeTotalCpuTimeNs”: 0, “segmentStatistics”: [], “traceInfo”: {}, “minConsumingFreshnessTimeMs”: 1645359464059, “numRowsResultSet”: 10
There is no ingestion right now, it’s stopped. Response is quite slow, in this case is event faster to use Trino and ORC
Few queries with really slow response here
r
50% of the time is spent decompressing raw string columns, 30% pruning segments via bloom filters, 15% looking up dictionary codes of string columns.
the 30% in segment pruning is caused by the high number of segments, improving this is something I have been thinking about recently
the 50% + 15% might be something you could avoid by querying for the columns you want rather than
select *
- do you actually want to select the raw string columns for whatever your use case is? This is all disk intensive...
so much of the profile is spent on selecting from string columns here that I would review data types (do they need to be strings?) storage (raw vs dictionary) - what is the length distribution of the raw string columns? It might make sense to use the V4 format for these columns (field level config - let me know if the raw strings have highly variable length or not) but also if
select *
is the right thing to be measuring
p
Thanks! Final projection will be much more narrow. But whats is weird is inconsistency it response time. Sometime I get timeout, sometime response in 3 seconds. Seems like there is some cache… but this could be related to disks as well. Till now all directs to SSD and conversion to numeric as much as possible 👍
Anyway the final outcome could be that i chose a wrong technology for my case, because I need to archive huge data with relatively low QPS max 10. And at the same time the scale of infra is low, max 10 servers. I'm considering Pinot or Clickhouse. But Pinot’s architecture is much better for us
k
Can you remove bloom filter?
The time does not make sense given the amount of work done from the response metadata
This query should be in milliseconds if you have the right layout