Skip to content

Commit 501a854

Browse files
committed
fix(data): fix DISTINCT deduplication to operate on projected output
SELECT DISTINCT was deduplicating on the raw document bytes before applying the projection. Two rows with the same projected value (e.g. same category) but different payload bytes were passed through as distinct, violating SQL semantics. Both the strict-schema and schemaless paths now project first, then deduplicate on the projected value.
1 parent 5293b97 commit 501a854

1 file changed

Lines changed: 51 additions & 35 deletions

File tree

  • nodedb/src/data/executor/handlers/document/read

nodedb/src/data/executor/handlers/document/read/scan.rs

Lines changed: 51 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -207,26 +207,30 @@ impl CoreLoop {
207207
if let Some(ref schema) = strict_schema
208208
&& window_specs.is_empty()
209209
{
210-
let deduped = if distinct {
211-
let mut seen = std::collections::HashSet::new();
212-
sorted
213-
.into_iter()
214-
.filter(|(_, value)| seen.insert(value.clone()))
215-
.collect::<Vec<_>>()
216-
} else {
217-
sorted
218-
};
219-
let result: Vec<_> = deduped
210+
// SQL DISTINCT semantics require deduplication on the
211+
// *projected* row, not the raw document bytes — two rows
212+
// with the same `category` but different ids/payload are
213+
// distinct as documents but the same under
214+
// `SELECT DISTINCT category`. Project first, then dedupe.
215+
let projected_rows: Vec<_> = sorted
220216
.into_iter()
221-
.skip(offset)
222-
.take(limit)
223217
.map(|(doc_id, val)| {
224218
let mp = decode_scanned_document_msgpack(&val, Some(schema));
225219
let projected =
226220
apply_projection_msgpack(&mp, &computed_cols, projection);
227221
(doc_id, projected)
228222
})
229223
.collect();
224+
let deduped = if distinct {
225+
let mut seen = std::collections::HashSet::new();
226+
projected_rows
227+
.into_iter()
228+
.filter(|(_, value)| seen.insert(value.clone()))
229+
.collect::<Vec<_>>()
230+
} else {
231+
projected_rows
232+
};
233+
let result: Vec<_> = deduped.into_iter().skip(offset).take(limit).collect();
230234
return self.send_document_rows_raw(task, &result, stream_chunk_size);
231235
}
232236

@@ -243,20 +247,10 @@ impl CoreLoop {
243247
&window_specs,
244248
);
245249

246-
let deduped: Vec<_> = if distinct {
247-
let mut seen = std::collections::HashSet::new();
248-
decoded_rows
249-
.into_iter()
250-
.filter(|(_, v)| seen.insert(v.to_string()))
251-
.collect()
252-
} else {
253-
decoded_rows
254-
};
255-
256-
let result: Vec<_> = deduped
250+
// Project first, then dedupe on the projected JSON value
251+
// so `SELECT DISTINCT col` honours SQL semantics.
252+
let projected_rows: Vec<_> = decoded_rows
257253
.into_iter()
258-
.skip(offset)
259-
.take(limit)
260254
.map(|(doc_id, data)| {
261255
let projected = apply_projection(data, &computed_cols, projection);
262256
DocumentRow {
@@ -266,35 +260,57 @@ impl CoreLoop {
266260
})
267261
.collect();
268262

269-
self.send_document_rows_transformed(task, &result, stream_chunk_size)
270-
} else {
271-
let deduped = if distinct {
263+
let deduped: Vec<_> = if distinct {
272264
let mut seen = std::collections::HashSet::new();
273-
sorted
265+
projected_rows
274266
.into_iter()
275-
.filter(|(_, value)| seen.insert(value.clone()))
267+
.filter(|row| seen.insert(row.data.to_string()))
276268
.collect()
277269
} else {
278-
sorted
270+
projected_rows
279271
};
280272

273+
let result: Vec<_> = deduped.into_iter().skip(offset).take(limit).collect();
274+
self.send_document_rows_transformed(task, &result, stream_chunk_size)
275+
} else {
281276
let needs_transform = !computed_cols.is_empty() || !projection.is_empty();
282277

283278
if needs_transform {
284-
let result: Vec<_> = deduped
279+
// Project first so DISTINCT acts on the projected
280+
// row, not the raw document.
281+
let projected_rows: Vec<_> = sorted
285282
.into_iter()
286-
.skip(offset)
287-
.take(limit)
288283
.map(|(doc_id, value)| {
289284
let mp = doc_format::json_to_msgpack(&value);
290285
let projected =
291286
apply_projection_msgpack(&mp, &computed_cols, projection);
292287
(doc_id, projected)
293288
})
294289
.collect();
295-
290+
let deduped = if distinct {
291+
let mut seen = std::collections::HashSet::new();
292+
projected_rows
293+
.into_iter()
294+
.filter(|(_, value)| seen.insert(value.clone()))
295+
.collect()
296+
} else {
297+
projected_rows
298+
};
299+
let result: Vec<_> = deduped.into_iter().skip(offset).take(limit).collect();
296300
self.send_document_rows_raw(task, &result, stream_chunk_size)
297301
} else {
302+
// No projection — `SELECT DISTINCT *` semantics dedupe
303+
// on the entire raw value, which is what the
304+
// pre-existing path does.
305+
let deduped = if distinct {
306+
let mut seen = std::collections::HashSet::new();
307+
sorted
308+
.into_iter()
309+
.filter(|(_, value)| seen.insert(value.clone()))
310+
.collect()
311+
} else {
312+
sorted
313+
};
298314
let rows: Vec<_> = deduped.into_iter().skip(offset).take(limit).collect();
299315
self.send_document_rows_raw(task, &rows, stream_chunk_size)
300316
}

0 commit comments

Comments
 (0)