If you are using Apache Pinot for telemetry, observability, or other high-volume event workloads, deletion may eventually become a requirement.
A customer might offboard and ask for their data to be removed. A bad instrumentation rollout might generate events you do not want to retain. Or a user might file a GDPR erasure request and their user id sits in a span attribute that was never supposed to be collected
A table may contain billions of events while the delete request might only cover a few hundred or a few thousand.
At first glance, that sounds like an easy task, but In Pinot, there’s more to it.
The reason comes down to how Pinot stores data: segments are written once and are not modified in place. To remove a row from a normal Pinot table, Pinot needs to create a new version of the segment without that row and replace the old segment.
To make this simpler, StarTree now offers a query based purge. You give StarTree’s Segment Purge Task a SQL WHERE clause to remove the matching rows from the live table while queries keep running.
But, before you run one, it helps to know how a purge works, how it is priced, and how to use it effectively.
Written once, read a billion times
If you’re new to Pinot, there are four components that do the heavy lifting.
- A segment is Pinot’s unit of storage. Think of it as a few million rows plus every index built over them, packed into one file.
- A minion is a worker node that rewrites those files in the background, so the servers answering your queries never stop to do maintenance.
- The controller decides what work exists and hands it out.
- A broker is the node a query lands on. It fans the query out to the servers holding the data and merges what comes back.
A segment is usually written once and never touched again. That sounds like a limitation, but it is the reason Pinot is fast. A file that never changes can be memory mapped, copied between replicas, cached aggressively and indexed without a single lock, because nothing has to check whether the bytes underneath it moved.
It also means there is no simple way to delete a row. What you do is build a new segment that is the old one minus the rows you do not want, then swap it in.
That is what a purge is. A minion downloads the segment, reads every row, drops the ones that match, writes a fresh segment out of what is left and uploads it. The controller swaps the new file in for the old one. From a query’s point of view the rows are simply gone.
So any way of deleting rows from Pinot comes down to two questions. Which segments get rewritten? And how does the rewrite know which rows to drop?
What open source Pinot gives you
Open source Pinot has several ways to remove data. Most of them work on whole segments rather than rows.
Retention is the cheapest delete there is. Once a segment’s time range falls outside the period the table keeps, the controller deletes the whole file. No rows are read and nothing is rebuilt. That is why nobody files a ticket for “everything older than thirty days”. Retention already took care of it.
The segment delete API removes whole segments on demand. That helps when the data you want gone happens to fill whole segments, like a bad backfill you pushed as its own batch. Telemetry deletes almost never line up that neatly.
Upsert tables have a delete column. Send a new record for a primary key with that column set to true and the key disappears from queries. It works well when it applies. You need an upsert table, you need to know every primary key you want gone and the old bytes stay on disk until a compaction task rewrites the segment. But most telemetry tables are append only, so this door is usually closed.
The last option is to rebuild the data yourself. Pull the affected time range from the source, filter out the rows, build new segments and push them with the same names so they replace the old ones. That works if you still have the source data and a batch pipeline. It also turns every delete request into a small project.
So open source Pinot can remove rows. What it cannot do is take the request the way it actually arrives, as a condition over a few columns on an append only table. It also cannot tell you what that will cost before it runs. That is the gap query based purge fills.
Write the delete the way the request arrives
StarTree’s Segment Purge Task started with two ways to choose rows. One is a file of keys, which fits a compliance workflow that hands over a list of user ids. The other is a small predicate language written as config keys. It has one comparison operator that applies to every field, one logical operator that joins them and a separate pair of keys for a time window.
That covered a lot of real requests. Purge these two tenants, purge this time range, purge this time range for these two tenants. Then we started getting asks that looked like this.
startTimeMs >= 1763078400000
AND startTimeMs <= 1763100000000
AND serviceName IN ('checkout-api', 'checkout-worker', ...six more)
AND spanKind NOT IN ('internal', 'producer', 'consumer')Code language: SQL (Structured Query Language) (sql)
That one condition has four different comparisons in it. The config language had room for one. IN was not on its list at all. Growing the config to fit would have meant building a small copy of SQL that nobody already knows, when Pinot already ships a SQL parser and a filter engine. So query based purge takes SQL.
Here is what that task config looks like.
{
"taskType": "SegmentPurgeTask",
"tableName": "myTable_REALTIME",
"taskConfigs": {
"purge.query.predicate": "WHERE \"startTimeMs\" >= 1763078400000 AND serviceName IN ('checkout-api', 'checkout-worker') AND spanKind NOT IN ('internal', 'producer')",
"purge.query.timeout.ms": "120000",
"purge.query.max-segments": "100000"
}
}Code language: JSON / JSON with Comments (json)
The config takes a WHERE clause and not a full SELECT. You own the condition, which is the part you care about. The task owns the rest of the query, so you never have to remember which hidden columns it needs.
Which SQL will a purge accept?
A purge predicate is not really a query. It is a test on a single row. Handed one row and nothing else, it has to answer keep or delete.
Anything that fits that test is fine. Comparisons, ranges and IN or NOT IN lists all work, joined with any mix of AND, OR and NOT. Anything that needs a second row or a second table to decide is rejected before the task starts.
- A JOIN would let a row in another table decide the fate of a row in this one. The minion rewriting this segment has never seen that other table.
- Aggregation, GROUP BY and HAVING make a row’s fate depend on the other rows in its group. Those rows sit in other segments on other machines. The minion doing the delete cannot see them.
- A subquery is the same problem in a different shape, with a hidden scan inside a predicate that is supposed to be cheap.
- DISTINCT shapes a result set. It does not test a row.
Those are the parts that are not supported. The validation error tells you which rule you hit. Because the config is a single WHERE expression, clauses like LIMIT or OFFSET cannot even be written into it. The parser rejects them first.
Will a purge slow down my dashboards?
A purge on a big table runs for a while. It runs on the same cluster that answers your customers’ queries. So the fair question is what it asks of your brokers and servers while it runs.
The obvious design has every minion ask a broker which rows in its segment match. That is one distributed query per segment, on the same brokers and thread pools as your customer traffic, for as long as the purge runs. On a table with tens of thousands of segments that adds up fast. It also has a subtle correctness gap. The broker answers from whatever copy of the segment the servers hold right now, while the minion is about to rewrite a copy it downloaded earlier.
Query based purge splits the job instead. Brokers answer one question, which segments hold matching rows. Minions answer the other, which rows inside a segment match, without asking anyone.
The first question goes to a broker as a single query. Pinot gives every table two hidden columns, $segmentName for the segment a row lives in and $docId for the row’s position inside it. The controller wraps your predicate like this.
SET useMultiStageEngine=true;
SELECT "$segmentName", COUNT(*)
FROM (SELECT "$segmentName", "$docId" FROM "myTable_REALTIME" WHERE <your predicate>) q
GROUP BY "$segmentName"
LIMIT <max segments + 1>Code language: SQL (Structured Query Language) (sql)
The GROUP BY matters more than it looks. It lets each server collapse its matches down to one row per segment before anything crosses the network. A SELECT DISTINCT gives the same answer but ships every matching row to the broker first. On a large table the DISTINCT form took over five minutes and the GROUP BY form came back in about forty seconds.
We also set a limit purge.query.max-segments, so the controller can tell when more segments matched than it is allowed to take. When that happens it purges the ones it found, logs a warning and picks the rest up on the next run. A delete too big for one pass takes several passes rather than refusing to start.
This listing query runs once per batch of subtasks. A batch is capped by table.max.num.tasks, which defaults to a thousand, so a purge that touches ten thousand segments sends about ten listing queries to your brokers and not ten thousand.
The second question never touches a broker. The minion downloads its segment, which it has to do anyway because rewriting it is the whole job. It loads the segment with the table’s index config, compiles your WHERE clause into a filter and runs that filter against the segment’s own indexes. What comes back is a bitmap of matching rows. The minion then walks the segment once and skips those rows as it writes the new one. Because the filter runs against the exact bytes being rewritten, there is no window for the answer to drift either.
The listing query is still a real query on your production brokers, so the task insists that it stays cheap. Every column in the WHERE clause needs an inverted, range or sorted index. The table’s time column is exempt, since segments already record the time range they cover. Bloom, JSON and text indexes do not count. If you know your table well enough to accept a scan, setting purge.query.force to true skips the check.
So how much does a delete cost?
A purge that removed a couple of million rows from a table holding several billion finished in about three minutes. That number tells you less than it seems to, because the thing that decided it is not the thing most people would guess.
A minion cannot edit a segment. To remove one row it downloads the whole file, reads every row, rebuilds every index over what is left, writes a new file and uploads it. That work is the same whether the predicate matched one row in the segment or all of them.
So the cost of a purge is set by how many segments your predicate lands in. The number of rows it matches barely matters.
Run that both ways before reading on. Ten million rows sitting in five segments is five rewrites. A hundred rows scattered one apiece across ten thousand segments is ten thousand rewrites, of ten thousand whole files, to remove a hundred rows. The second delete is a hundred thousand times smaller and two thousand times more work.
That flips the instinct most people bring from a row store, where deleting more rows costs more. Here, deleting more rows is often free. Deleting them from more places is what costs.
The lever you have is the time filter. On a telemetry table it is a big one. Segments are written in the order events arrive, so each one covers a narrow slice of time and records exactly which slice. Put a time range in the predicate and the listing query lands on a narrow band of segments. Leave it out and the same rows are spread across every segment in your retention window.
Some deletes honestly have no time bound. An offboarded tenant’s spans really are everywhere. Plenty of others do have one, though. A bad rollout has a start and an end. People leave it out of the predicate because it feels redundant next to the service name. It is not redundant. It is the difference between a hundred rewrites and ten thousand.
Two more costs are worth planning for on a big purge. Each touched segment becomes its own subtask, so ten thousand segments means ten thousand downloads, rebuilds and uploads, spread over batches. And the trigger is synchronous. Both the dry run and the execute call run the listing query before they return, so on a very large table the HTTP call can sit there for a couple of minutes. The default purge.query.timeout.ms is thirty seconds, which is not enough for a listing query over tens of thousands of segments, so raising it is usually the first change you make.
Look before you delete
A delete has no undo, so the dry run is not optional advice. It is the procedure.
A dry run takes the same task configs as the real run and posts them to one endpoint.
POST /tasks/SegmentPurgeTask/{tableNameWithType}/dryRun?verbose=trueCode language: JavaScript (javascript)
It runs the same validation and the same listing query, plus a COUNT(*) around your predicate. Then it stops. Nothing is downloaded and nothing is rewritten. What comes back answers four questions.
Is the predicate right? The matching row count is the number to hold against what you expected. A predicate that matches four hundred times more rows than you had in mind is the most common thing a dry run catches.
How much work is this? The number of segments that would be rewritten. After the last section, you know this is the real price of the job.
How long will it take? The number of subtasks, plus how many batches that becomes at your table.max.num.tasks. That turns “this seems big” into a number you can put in a change ticket.
What is being left out? Segments the purge passed over, like ones still consuming, ones with no records and ones already handled in this purge cycle. With verbose set you get them grouped by reason and by name, which is what you want when the question is why it skipped the one segment you care about.
On top of those it raises warnings, not errors, when the delete looks risky. It warns when the purge touches ten thousand segments or more, with a stronger warning when that is more than half the table. It also warns when an earlier purge on the table is still incomplete and when the predicate matches nothing. On an upsert table it warns if refresh, compaction or snapshot tasks are running. None of those block you, because deleting most of a table is sometimes exactly the request. They make you say it out loud first.
Validation errors show up here too. An unindexed filter column, a full SELECT instead of a WHERE clause and an upsert table with no delete column all fail in the dry run rather than halfway through a thousand subtasks.
Two honest caveats. A dry run executes the real listing query, so on a very large table it takes the same couple of minutes the real selection does. And every number in it is a snapshot. Rows keep arriving, so what you saw is what was true when you asked.
Deleting from a table that keeps every version
Upsert tables work differently. They are also where people get caught out the most.
An upsert table keeps every version of a row it ever ingested. Pinot keeps a live view on top that says, for each primary key, which single row is the current one. Every query reads through that view, so you only see the winner.
The purge uses the delete column that upsert tables already have, rather than inventing a second way to delete. It does not drop the matching rows. It sets the delete column to true on them, in place, in the segment where they already sit. Everything else in the row stays the same, including the value that decides which version wins.
Two things follow from that.
The table does not get smaller. Same rows in, same rows out, one column different. The rows vanish from query results as soon as the new segment is swapped in, but the bytes only go away when compaction reaches them. So validation refuses a purge on an upsert table unless it has a delete column and a SegmentRefreshTask or UpsertCompactionTask configured to reclaim the space.
Do not run it next to tasks that rewrite the same segments. Refresh, compaction and snapshot tasks all rewrite upsert segments too. The dry run warns when any of them are active, so check it before you trigger.
Could it delete the wrong rows?
A feature that deletes production data from a string in a config file needs guardrails that protect the rows you meant to keep.
The segment is checked twice. Before a minion starts, it checks that the segment is still the version the purge was planned against. If it changed, the minion skips the work. The upload then carries the segment’s original checksum as a precondition. If something else rewrote the segment in the meantime, the controller rejects the upload instead of letting the purge overwrite newer data.
Rows are not sent through ingestion a second time. The code that rebuilds a segment during a purge is the same code that builds one during ingestion. By default that code applies the table’s ingestion transforms and filters. Every row in a segment already went through those once on the way in. Running them again would compute derived columns from values that were already derived. Worse, an ingestion filter could drop rows nobody asked to purge while the task still reports success. So on current releases the purge rebuild turns ingestion transforms and filters off by default, with per task settings to turn individual ones back on when a table genuinely needs them.
One way of choosing rows at a time. The key file, the config predicate and the SQL predicate cannot be mixed in one task. Trying is a validation error. Changing the predicate also resets the purge’s progress, so a half finished purge under the old condition never blends into the new one.
Is the data really gone?
For a GDPR request, “gone from query results” and “gone from storage” are different promises, so it is worth knowing where the bytes end up.
On an append only table, the rewritten segment replaces the old one on every replica. It also overwrites the old file in deep store under the same name. If every row in a segment matched, the minion deletes the segment instead of uploading an empty one. Its file then follows Pinot’s normal deleted segment path. It is kept briefly as a replaced segment, then moved to a deleted segments folder that holds files for seven days by default. You can shorten that window per table with deletedSegmentsRetentionPeriod if your compliance clock runs faster.
On an upsert table the rows disappear from queries as soon as the purge lands. The bytes leave when the next refresh or compaction rewrites that segment.
And Pinot can only purge what Pinot holds. The Kafka topic or the files you ingested from keep their own copy until their own retention expires.
Before your next delete request
When the next request lands, write it as a WHERE clause and add the time range whenever the request has one. Dry run it before anything ejlse. Check the matching row count against what you expected, then read the segment count as the bill, because that is the number that decides how long the purge runs. Raise purge.query.timeout.ms for big tables. On an upsert table, make sure refresh or compaction is running, because that is when the bytes actually leave.
Query based purge ships in StarTree Cloud 0.14.0 and later. The Segment Purge Task docs cover the full config, the dry run API and every warning it can raise. If you already run StarTree, the quickest way to see what a delete would cost is to dry run one against your own table, since a dry run never deletes anything.

