DEV Community

Franck Pachot
Franck Pachot

Posted on

DocumentDB 0.116: $sort $group prefix pushdown

DocumentDB 0.116-0, released August 20, 2026, introduced prefix sort pushdown for $group accumulators. It lets the grouping stage consume rows in the existing sort order, avoiding a separate blocking sort or hash aggregate.

The previous article in this series covered a different optimization: a distinct scan that skips duplicate index entries when a group returns one row per distinct key. Prefix sort pushdown does not skip entries. It removes the top-level sort. The query pattern is also related to the first-per-group example in First/Last per Group: PostgreSQL DISTINCT ON and MongoDB DISTINCT_SCAN Performance, which shows MongoDB using DISTINCT_SCAN.

DocumentDB implements the MongoDB API as a fully open-source PostgreSQL extension. It offers an open alternative for MongoDB applications and exposes aggregation execution through both MongoDB-like and PostgreSQL plans, providing insight into sorts, buffers, and temporary I/O. Microsoft is the main contributor, with improvements coming from enterprise experience with Azure DocumentDB.

To demonstrate prefix sort pushdown, I load 50,000 documents across 100 groups, add a 200-byte payload, and create an ordered index on {a: 1}. The pipeline sorts by a, groups by the same key, and returns the first name in each group. Because the group key is a prefix of the sort keys, the accumulator can use the existing index order instead of sorting all 50,000 rows again.

Controlled method

MongoDB, DocumentDB 0.114, and DocumentDB 0.116 receive the same data, ordered index, hint, and pipeline. Both DocumentDB versions get the same VACUUM (ANALYZE), and both native runs use the same planner settings. The documentdb.enableSortPushToAccumulatorWithPrefix setting exists in 0.114 but is disabled by default. It is enabled by default in 0.116. I compare explicit sort nodes, temporary I/O, rows, and buffers rather than total elapsed time. The only timing I cite is when the gateway's group stage starts producing rows (executionStartAtTimeMillis), which shows whether grouping blocks or streams. In the native (SQL function) runs, I set enable_hashagg=off (and enable_seqscan / enable_bitmapscan) to force the sorted plan, but the gateway (MongoDB-compatible endpoint) runs use default settings.

MongoDB 8.0 reference

I first run the pipeline on MongoDB Atlas 8.0:

db = db.getSiblingDB("perf116");
db.sort_group.drop();
db.sort_group.insertMany(Array.from({length: 50000}, (_, index) => ({
  _id: index + 1,
  a: (index + 1) % 100,
  b: (50000 - index - 1) % 1000,
  name: `name_${index + 1}`,
  payload: "x".repeat(200)
})));
db.sort_group.createIndex({a: 1});

const pipeline = [
  {$sort: {a: 1}},
  {$group: {_id: "$a", firstVal: {$first: "$name"}}}
];
print(EJSON.stringify({
  resultCount: db.sort_group.aggregate(
    pipeline,
    {hint: "a_1"}
  ).toArray().length,
  explain: db.sort_group.explain("executionStats").aggregate(
    pipeline,
    {hint: "a_1"}
  )
}, null, 2));
Enter fullscreen mode Exit fullscreen mode

It shows a DISTINCT_SCAN:

{
  "resultCount": 100,
  "explain": {
    "explainVersion": "1",
    "stages": [
      {
        "$cursor": {
          "queryPlanner": {
            "namespace": "perf116.sort_group",
            "parsedQuery": {},
            "indexFilterSet": false,
            "queryHash": "DC5C6196",
            "planCacheShapeHash": "DC5C6196",
            "planCacheKey": "979F6906",
            "optimizationTimeMillis": 0,
            "maxIndexedOrSolutionsReached": false,
            "maxIndexedAndSolutionsReached": false,
            "maxScansToExplodeReached": false,
            "prunedSimilarIndexes": false,
            "winningPlan": {
              "isCached": false,
              "stage": "FETCH",
              "inputStage": {
                "stage": "DISTINCT_SCAN",
                "keyPattern": {
                  "a": 1
                },
                "indexName": "a_1",
                "isMultiKey": false,
                "multiKeyPaths": {
                  "a": []
                },
                "isUnique": false,
                "isSparse": false,
                "isPartial": false,
                "indexVersion": 2,
                "direction": "forward",
                "indexBounds": {
                  "a": [
                    "[MinKey, MaxKey]"
                  ]
                }
              }
            },
            "rejectedPlans": []
          },
          "executionStats": {
            "executionSuccess": true,
            "nReturned": 100,
            "executionTimeMillis": 2,
            "totalKeysExamined": 100,
            "totalDocsExamined": 100,
            "executionStages": {
              "isCached": false,
              "stage": "FETCH",
              "nReturned": 100,
              "executionTimeMillisEstimate": 0,
              "works": 101,
              "advanced": 100,
              "needTime": 0,
              "needYield": 0,
              "saveState": 3,
              "restoreState": 3,
              "isEOF": 1,
              "docsExamined": 100,
              "alreadyHasObj": 0,
              "inputStage": {
                "stage": "DISTINCT_SCAN",
                "nReturned": 100,
                "executionTimeMillisEstimate": 0,
                "works": 101,
                "advanced": 100,
                "needTime": 0,
                "needYield": 0,
                "saveState": 3,
                "restoreState": 3,
                "isEOF": 1,
                "keyPattern": {
                  "a": 1
                },
                "indexName": "a_1",
                "isMultiKey": false,
                "multiKeyPaths": {
                  "a": []
                },
                "isUnique": false,
                "isSparse": false,
                "isPartial": false,
                "indexVersion": 2,
                "direction": "forward",
                "indexBounds": {
                  "a": [
                    "[MinKey, MaxKey]"
                  ]
                },
                "keysExamined": 100
              }
            }
          }
        },
        "nReturned": 100,
        "executionTimeMillisEstimate": 3
      },
      {
        "$groupByDistinctScan": {
          "newRoot": {
            "_id": "$a",
            "firstVal": "$name"
          }
        },
        "nReturned": 100,
        "executionTimeMillisEstimate": 3
      }
    ],
    "queryShapeHash": "0822C17532FC5F85649D595163EF78A0BAEF5CDA8811F1689D203E1AD79FD2F7",
    "serverInfo": {
      "host": "a004f7434a57",
      "port": 27017,
      "version": "8.0.28",
      "gitVersion": "cd6fc9b3b7cf87ff2bbca0af67382ac407fc682a"
    },
    "serverParameters": {
      "internalQueryFacetBufferSizeBytes": 104857600,
      "internalQueryFacetMaxOutputDocSizeBytes": 104857600,
      "internalLookupStageIntermediateDocumentMaxSizeBytes": 104857600,
      "internalDocumentSourceGroupMaxMemoryBytes": 104857600,
      "internalQueryMaxBlockingSortMemoryUsageBytes": 104857600,
      "internalQueryProhibitBlockingMergeOnMongoS": 0,
      "internalQueryMaxAddToSetBytes": 104857600,
      "internalDocumentSourceSetWindowFieldsMaxMemoryBytes": 104857600,
      "internalQueryFrameworkControl": "trySbeRestricted",
      "internalQueryPlannerIgnoreIndexWithCollationForRegex": 1
    },
    "command": {
      "aggregate": "sort_group",
      "pipeline": [
        {
          "$sort": {
            "a": 1
          }
        },
        {
          "$group": {
            "_id": "$a",
            "firstVal": {
              "$first": "$name"
            }
          }
        }
      ],
      "hint": "a_1",
      "cursor": {},
      "$db": "perf116"
    },
    "ok": 1
  }
}
Enter fullscreen mode Exit fullscreen mode

MongoDB recognizes the first-per-group pattern. It uses DISTINCT_SCAN on a_1, examines 100 index keys, fetches 100 documents for name, and returns 100 groups. There is no blocking sort.

DocumentDB 0.114 (before this optimization)

I load and index the same data through the DocumentDB 0.114 gateway (MongoDB-compatible endpoint):

db = db.getSiblingDB("perf116");
db.sort_group.drop();
db.sort_group.insertMany(Array.from({length: 50000}, (_, index) => ({
  _id: index + 1,
  a: (index + 1) % 100,
  b: (50000 - index - 1) % 1000,
  name: `name_${index + 1}`,
  payload: "x".repeat(200)
})));
db.sort_group.createIndex(
  {a: 1},
  {storageEngine: {enableOrderedIndex: true}}
);
print(EJSON.stringify({
  insertedDocuments: db.sort_group.countDocuments(),
  indexes: db.sort_group.getIndexes().map(index => index.name)
}, null, 2));


{
  "insertedDocuments": 50000,
  "indexes": [
    "_id_",
    "a_1"
  ]
}
Enter fullscreen mode Exit fullscreen mode

Connected to the PostgreSQL endpoint, I make the table statistics and visibility state deterministic:

\pset pager off
SET search_path TO documentdb_api_catalog, public;
SELECT extversion AS documentdb_version
FROM pg_extension
WHERE extname = 'documentdb';
SELECT collection_id
FROM collections
WHERE database_name = 'perf116' AND collection_name = 'sort_group'
\gset
VACUUM (ANALYZE) documentdb_data.documents_:collection_id;
SELECT :'collection_id' AS vacuumed_collection_id;

 documentdb_version 
--------------------
 0.114-0
(1 row)

 vacuumed_collection_id 
------------------------
 2
(1 row)
Enter fullscreen mode Exit fullscreen mode

Then I run the same hinted pipeline on the MongoSH:

db = db.getSiblingDB("perf116");
const pipeline = [
  {$sort: {a: 1}},
  {$group: {_id: "$a", firstVal: {$first: "$name"}}}
];
print(EJSON.stringify({
  resultCount: db.sort_group.aggregate(
    pipeline,
    {hint: "a_1"}
  ).toArray().length,
  explain: db.sort_group.explain("executionStats").aggregate(
    pipeline,
    {hint: "a_1"}
  )
}, null, 2));


{
  "resultCount": 100,
  "explain": {
    "explainVersion": 2,
    "command": "db.runCommand({explain: { 'aggregate': 'sort_group', 'pipeline': [{ '$sort': { 'a': 1 } }, { '$group': { '_id': '$a', 'firstVal': { '$first': '$name' } } }], 'hint': 'a_1', 'cursor': {} }})",
    "explainCommandPlanningTimeMillis": 2.649,
    "explainCommandExecTimeMillis": 349.424,
    "stages": [
      {
        "$cursor": {
          "queryPlanner": {
            "winningPlan": {
              "stage": "FETCH",
              "startupCost": 0,
              "totalCost": 13.89,
              "estimatedTotalKeysExamined": 5556,
              "inputStage": {
                "stage": "IXSCAN",
                "indexName": "a_1",
                "direction": "Forward",
                "startupCost": 0,
                "totalCost": 13.89,
                "hasOrderBy": true,
                "indexFilterSet": [
                  {
                    "a": {
                      "$range": {
                        "orderByScan": 1
                      }
                    }
                  }
                ],
                "estimatedTotalKeysExamined": 5556
              }
            }
          },
          "executionStats": {
            "nReturned": 50000,
            "executionTimeMillis": 125.407,
            "executionStartAtTimeMillis": 0.036,
            "totalDocsExamined": 50000,
            "totalKeysExamined": 50000,
            "executionStages": {
              "stage": "FETCH",
              "nReturned": 50000,
              "executionTimeMillis": 125.407,
              "executionStartAtTimeMillis": 0.036,
              "totalKeysExamined": 50000,
              "numBlocksFromCache": 50023,
              "inputStage": {
                "stage": "IXSCAN",
                "nReturned": 50000,
                "executionTimeMillis": 125.407,
                "executionStartAtTimeMillis": 0.036,
                "indexName": "a_1",
                "totalKeysExamined": 50000,
                "numBlocksFromCache": 50023
              }
            }
          }
        }
      },
      {
        "$group": {
          "queryPlanner": {
            "winningPlan": {
              "stage": "GROUP",
              "startupCost": 111.12,
              "totalCost": 236.13,
              "aggStrategy": "Hashed",
              "estimatedTotalKeysExamined": 5556,
              "inputStage": {
                "stage": "PROJECTION_DEFAULT",
                "startupCost": 0,
                "totalCost": 83.34,
                "estimatedTotalKeysExamined": 5556
              }
            }
          },
          "executionStats": {
            "nReturned": 100,
            "executionTimeMillis": 349.126,
            "executionStartAtTimeMillis": 348.869,
            "totalDocsExamined": 100,
            "totalKeysExamined": 100,
            "executionStages": {
              "stage": "GROUP",
              "nReturned": 100,
              "executionTimeMillis": 349.126,
              "executionStartAtTimeMillis": 348.869,
              "totalDocsExamined": 100,
              "totalKeysExamined": 100,
              "numBlocksFromCache": 50023,
              "inputStage": {
                "stage": "PROJECTION_DEFAULT",
                "nReturned": 50000,
                "executionTimeMillis": 255.141,
                "executionStartAtTimeMillis": 0.041,
                "totalDocsExamined": 50000,
                "totalKeysExamined": 50000,
                "numBlocksFromCache": 50023
              }
            }
          }
        }
      }
    ],
    "ok": 1
  }
}
Enter fullscreen mode Exit fullscreen mode

With default settings, the gateway plan uses a hash aggregate (aggStrategy: "Hashed"), which cannot return any group until it reads all 50,000 rows. The group stage starts at 348.9 ms. The MongoDB-compatible explain doesn't show buffers or temporary I/O, so I inspect the native plan from the PostgreSQL endpoint. I disable hash aggregation there to see the sort-based plan that the new optimization targets:

\pset pager off
\set ON_ERROR_STOP on
SET search_path TO documentdb_api, documentdb_core, documentdb_api_catalog,
  documentdb_api_internal, public;

SELECT extversion AS documentdb_version
FROM pg_extension
WHERE extname = 'documentdb';

SELECT documentdb_api.drop_collection('perf116', 'sort_group');
SELECT count(documentdb_api.insert_one(
  'perf116',
  'sort_group',
  format(
    '{"_id":%s,"a":%s,"b":%s,"name":"name_%s","payload":"%s"}',
    i, i % 100, (50000 - i) % 1000, i, repeat('x', 200)
  )::documentdb_core.bson,
  NULL
))
FROM generate_series(1, 50000) i;

SELECT documentdb_api_internal.create_indexes_non_concurrently(
  'perf116',
  '{"createIndexes":"sort_group","indexes":[{"key":{"a":1},"storageEngine":{"enableOrderedIndex":true},"name":"a_1"}]}',
  true
);

SELECT collection_id
FROM documentdb_api_catalog.collections
WHERE database_name = 'perf116' AND collection_name = 'sort_group'
\gset
VACUUM (ANALYZE) documentdb_data.documents_:collection_id;

SET enable_seqscan TO off;
SET enable_bitmapscan TO off;
SET enable_hashagg TO off;
SET documentdb.enableGroupByCompoundIdIndexPushdown TO off;

\echo === DocumentDB 0.114 ===
EXPLAIN (ANALYZE, VERBOSE, COSTS OFF, SUMMARY OFF, TIMING OFF, BUFFERS)
SELECT document
FROM bson_aggregation_pipeline(
  'perf116',
  '{"aggregate":"sort_group","pipeline":[{"$sort":{"a":1}},{"$group":{"_id":"$a","firstVal":{"$first":"$name"}}}],"cursor":{},"hint":"a_1"}'
);
Enter fullscreen mode Exit fullscreen mode

Here is the output:

 documentdb_version 
--------------------
 0.114-0
(1 row)

 drop_collection 
-----------------
 f
(1 row)

NOTICE:  creating collection
 count 
-------
 50000
(1 row)

                                                                                                                   create_indexes_non_concurrently                                                                                                                   
---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 BSONHEX7e00000003726177006c0000000364656661756c7453686172640059000000106e756d496e64657865734265666f72650001000000106e756d496e6465786573416674657200020000000863726561746564436f6c6c656374696f6e4175746f6d61746963616c6c790000106f6b00010000000000106f6b000100000000
(1 row)

=== DocumentDB 0.114 ===
                                                                                                                                                                                                                                                                                        QUERY PLAN                                                                                                                                                                                                                                                                                         
-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 GroupAggregate (actual rows=100 loops=1)
   Output: bson_repath_and_build('_id'::text, (bson_expression_get(agg_stage_0.document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700210816cda001000000'::bson)), 'firstVal'::text, bson_expression_get(bsonfirst(agg_stage_0.document, '{BSONHEX0c0000001061000100000000}'::bson[]), 'BSONHEX11000000020006000000246e616d650000'::bson, true, 'BSONHEX12000000096e6f7700210816cda001000000'::bson)), (bson_expression_get(agg_stage_0.document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700210816cda001000000'::bson))
   Group Key: (bson_expression_get(agg_stage_0.document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700210816cda001000000'::bson))
   Buffers: shared hit=50023, temp read=1768 written=1771
   ->  Sort (actual rows=50000 loops=1)
         Output: (bson_expression_get(agg_stage_0.document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700210816cda001000000'::bson)), agg_stage_0.document
         Sort Key: (bson_expression_get(agg_stage_0.document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700210816cda001000000'::bson))
         Sort Method: external merge  Disk: 14144kB
         Buffers: shared hit=50023, temp read=1768 written=1771
         ->  Subquery Scan on agg_stage_0 (actual rows=50000 loops=1)
               Output: bson_expression_get(agg_stage_0.document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700210816cda001000000'::bson), agg_stage_0.document
               Buffers: shared hit=50023
               ->  Index Scan using a_1 on documentdb_data.documents_2 collection (actual rows=50000 loops=1)
                     Output: collection.document, bson_orderby(collection.document, 'BSONHEX0c0000001061000100000000'::bson)
                     Index Cond: (collection.document @<> 'BSONHEX1e00000003610016000000106f7264657242795363616e00010000000000'::bson)
                     Order By: (collection.document |-<> 'BSONHEX0c0000001061000100000000'::bson)
                     Buffers: shared hit=50023
 Planning:
   Buffers: shared hit=548
(19 rows)
Enter fullscreen mode Exit fullscreen mode

With hash aggregation disabled, DocumentDB 0.114 reads the ordered a_1 index but still sorts all 50,000 rows before GroupAggregate. The sort spills to disk as a 14,144 kB external merge and reports 1,768 temporary blocks read and 1,771 written.

DocumentDB 0.116 (after this optimization)

I repeat the same:

db = db.getSiblingDB("perf116");
db.sort_group.drop();
db.sort_group.insertMany(Array.from({length: 50000}, (_, index) => ({
  _id: index + 1,
  a: (index + 1) % 100,
  b: (50000 - index - 1) % 1000,
  name: `name_${index + 1}`,
  payload: "x".repeat(200)
})));
db.sort_group.createIndex(
  {a: 1},
  {storageEngine: {enableOrderedIndex: true}}
);
print(EJSON.stringify({
  insertedDocuments: db.sort_group.countDocuments(),
  indexes: db.sort_group.getIndexes().map(index => index.name)
}, null, 2));


{
  "insertedDocuments": 50000,
  "indexes": [
    "_id_",
    "a_1"
  ]
}
Enter fullscreen mode Exit fullscreen mode

The same maintenance step is applied:

\pset pager off
SET search_path TO documentdb_api_catalog, public;
SELECT extversion AS documentdb_version
FROM pg_extension
WHERE extname = 'documentdb';
SELECT collection_id
FROM collections
WHERE database_name = 'perf116' AND collection_name = 'sort_group'
\gset
VACUUM (ANALYZE) documentdb_data.documents_:collection_id;
SELECT :'collection_id' AS vacuumed_collection_id;


 documentdb_version 
--------------------
 0.116-0
(1 row)

 vacuumed_collection_id 
------------------------
 7
(1 row)
Enter fullscreen mode Exit fullscreen mode

The query and hint are unchanged:

db = db.getSiblingDB("perf116");
const pipeline = [
  {$sort: {a: 1}},
  {$group: {_id: "$a", firstVal: {$first: "$name"}}}
];
print(EJSON.stringify({
  resultCount: db.sort_group.aggregate(
    pipeline,
    {hint: "a_1"}
  ).toArray().length,
  explain: db.sort_group.explain("executionStats").aggregate(
    pipeline,
    {hint: "a_1"}
  )
}, null, 2));


{
  "resultCount": 100,
  "explain": {
    "explainVersion": 2,
    "command": "db.runCommand({explain: { 'aggregate': 'sort_group', 'pipeline': [{ '$sort': { 'a': 1 } }, { '$group': { '_id': '$a', 'firstVal': { '$first': '$name' } } }], 'hint': 'a_1', 'cursor': {} }})",
    "explainCommandPlanningTimeMillis": 1.094,
    "explainCommandExecTimeMillis": 280.012,
    "stages": [
      {
        "$cursor": {
          "queryPlanner": {
            "winningPlan": {
              "stage": "FETCH",
              "ns": "perf116.sort_group",
              "startupCost": 0,
              "totalCost": 0,
              "estimatedTotalKeysExamined": 5556,
              "inputStage": {
                "stage": "IXSCAN",
                "ns": "perf116.sort_group",
                "indexName": "a_1",
                "direction": "Forward",
                "indexUsage": {
                  "indexKeyString": "{\"a\": 1}",
                  "isMultiKey": false,
                  "bounds": [
                    "[\"a\": (MinKey, MaxKey)]"
                  ]
                },
                "startupCost": 0,
                "totalCost": 0,
                "hasOrderBy": true,
                "indexFilterSet": [
                  {
                    "a": {
                      "$range": {
                        "orderByScan": 1
                      }
                    }
                  }
                ],
                "estimatedTotalKeysExamined": 5556
              }
            },
            "indexCosts": [
              {
                "namespace": "perf116.sort_group",
                "costs": [
                  {
                    "indexName": "_id_",
                    "startupCost": 0.415,
                    "totalCost": 1379.415,
                    "selectivity": 1,
                    "correlation": 0.75,
                    "estimatedPercentIndexPagesLoaded": 100,
                    "estimatedTotalIndexEntries": 50000,
                    "boundarySelectivity": 1
                  }
                ]
              }
            ]
          },
          "executionStats": {
            "nReturned": 50000,
            "executionTimeMillis": 92.674,
            "executionStartAtTimeMillis": 0.11,
            "totalDocsExamined": 50000,
            "totalKeysExamined": 50000,
            "executionStages": {
              "stage": "FETCH",
              "nReturned": 50000,
              "executionTimeMillis": 92.674,
              "executionStartAtTimeMillis": 0.11,
              "totalKeysExamined": 50000,
              "numBlocksFromCache": 50044,
              "inputStage": {
                "stage": "IXSCAN",
                "nReturned": 50000,
                "executionTimeMillis": 92.674,
                "executionStartAtTimeMillis": 0.11,
                "indexName": "a_1",
                "indexUsage": {
                  "scanLoops": 100,
                  "scanType": "ordered",
                  "scanKeys": [
                    "key 1: [(isInequality: true, estimatedEntryCount: 50000)]"
                  ]
                },
                "totalKeysExamined": 50000,
                "numBlocksFromCache": 50044
              }
            }
          }
        }
      },
      {
        "$group": {
          "queryPlanner": {
            "winningPlan": {
              "stage": "GROUP",
              "startupCost": 0,
              "totalCost": 152.79,
              "aggStrategy": "Sorted",
              "estimatedTotalKeysExamined": 5556
            }
          },
          "executionStats": {
            "nReturned": 100,
            "executionTimeMillis": 279.817,
            "executionStartAtTimeMillis": 3.311,
            "totalDocsExamined": 100,
            "totalKeysExamined": 100,
            "executionStages": {
              "stage": "GROUP",
              "nReturned": 100,
              "executionTimeMillis": 279.817,
              "executionStartAtTimeMillis": 3.311,
              "totalDocsExamined": 100,
              "totalKeysExamined": 100,
              "numBlocksFromCache": 50044
            }
          }
        }
      }
    ],
    "ok": 1
  }
}
Enter fullscreen mode Exit fullscreen mode

The group stage now reports aggStrategy: "Sorted" instead of 0.114's "Hashed": with the $sort pushed down, $group consumes the ordered index scan directly.

The gateway still reports 50,000 index entries and documents. DocumentDB's distinct scan for $first-only groups (documentdb.enableDistinctScanForGroupFirst, a sibling of the enableGroupByDistinctScan setting from the previous article) is disabled by default in 0.116, so all 50,000 keys are examined. I left it disabled to isolate the sort elimination, which is visible in the PostgreSQL execution plan (you can try all variations in this db<>fiddle):

\pset pager off
\set ON_ERROR_STOP on
SET search_path TO documentdb_api, documentdb_core, documentdb_api_catalog,
  documentdb_api_internal, public;

SELECT extversion AS documentdb_version
FROM pg_extension
WHERE extname = 'documentdb';

SELECT documentdb_api.drop_collection('perf116', 'sort_group');
SELECT count(documentdb_api.insert_one(
  'perf116',
  'sort_group',
  format(
    '{"_id":%s,"a":%s,"b":%s,"name":"name_%s","payload":"%s"}',
    i, i % 100, (50000 - i) % 1000, i, repeat('x', 200)
  )::documentdb_core.bson,
  NULL
))
FROM generate_series(1, 50000) i;

SELECT documentdb_api_internal.create_indexes_non_concurrently(
  'perf116',
  '{"createIndexes":"sort_group","indexes":[{"key":{"a":1},"storageEngine":{"enableOrderedIndex":true},"name":"a_1"}]}',
  true
);

SELECT collection_id
FROM documentdb_api_catalog.collections
WHERE database_name = 'perf116' AND collection_name = 'sort_group'
\gset
VACUUM (ANALYZE) documentdb_data.documents_:collection_id;

SET enable_seqscan TO off;
SET enable_bitmapscan TO off;
SET enable_hashagg TO off;
SET documentdb.enableGroupByCompoundIdIndexPushdown TO off;

SET documentdb.enableSortPushToAccumulatorWithPrefix TO on;
\echo === DocumentDB 0.116 feature ON ===
EXPLAIN (ANALYZE, VERBOSE, COSTS OFF, SUMMARY OFF, TIMING OFF, BUFFERS)
SELECT document
FROM bson_aggregation_pipeline(
  'perf116',
  '{"aggregate":"sort_group","pipeline":[{"$sort":{"a":1}},{"$group":{"_id":"$a","firstVal":{"$first":"$name"}}}],"cursor":{},"hint":"a_1"}'
);
Enter fullscreen mode Exit fullscreen mode

Here is the output:

 documentdb_version 
--------------------
 0.116-0
(1 row)

 drop_collection 
-----------------
 t
(1 row)

NOTICE:  creating collection
 count 
-------
 50000
(1 row)

                                                                                                                   create_indexes_non_concurrently                                                                                                                   
---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 BSONHEX7e00000003726177006c0000000364656661756c7453686172640059000000106e756d496e64657865734265666f72650001000000106e756d496e6465786573416674657200020000000863726561746564436f6c6c656374696f6e4175746f6d61746963616c6c790000106f6b00010000000000106f6b000100000000
(1 row)

=== DocumentDB 0.116 feature ON ===
                                                                                                                                                                                                                                            QUERY PLAN                                                                                                                                                                                                                                             
---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 GroupAggregate (actual rows=100 loops=1)
   Output: bson_repath_and_build('_id'::text, (bson_expression_get(document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700d2f01acda001000000'::bson)), 'firstVal'::text, bsonfirstwithexpr(document, 'BSONHEX11000000020006000000246e616d650000'::bson, 'BSONHEX12000000096e6f7700d2f01acda001000000'::bson, NULL::text)), (bson_expression_get(document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700d2f01acda001000000'::bson))
   Group Key: bson_expression_get(collection.document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700d2f01acda001000000'::bson)
   Buffers: shared hit=50023
   ->  Custom Scan (DocumentDBApiExplainQueryScan) (actual rows=50000 loops=1)
         Output: bson_expression_get(document, 'BSONHEX0e00000002000300000024610000'::bson, true, 'BSONHEX12000000096e6f7700d2f01acda001000000'::bson), document
         namespaceName: perf116.sort_group
         indexName: a_1
         indexKey: {"a": 1}
         isMultiKey: false
         indexBounds: ["a": (MinKey, MaxKey)]
         innerScanLoops: 100 loops
         scanType: ordered
         scanKeyDetails: key 1: [(isInequality: true, estimatedEntryCount: 50000)]
         _id_: (startup cost=0.415, total cost=1379.415, selectivity=1, correlation=0.750, estimated index pages loaded=100.00%, estimated total index entries=50000, boundary selectivity=1, num boundaries=0, estimated data pages loaded=0.00%)
         Buffers: shared hit=50023
         ->  Index Scan using a_1 on documentdb_data.documents_8 collection (actual rows=50000 loops=1)
               Output: document
               Index Cond: (collection.document @<> 'BSONHEX1e00000003610016000000106f7264657242795363616e00010000000000'::bson)
               Order By: (collection.document |-<> 'BSONHEX0c0000001061000100000000'::bson)
               Buffers: shared hit=50023
 Planning:
   Buffers: shared hit=387
(23 rows)
Enter fullscreen mode Exit fullscreen mode

In 0.116, GroupAggregate consumes the ordered index scan directly. The Sort node, 14 MB spill, and all temporary reads and writes are gone. The scan work is unchanged, which isolates the improvement to sort elimination.

Conclusion

DocumentDB 0.116 turns a blocking $group into a streaming one: groups come directly from the ordered index scan, with no hash table or sort in between.

Here is a summary of the experiment:

Engine and API Index entries Documents fetched Blocking step Temporary I/O
MongoDB 8.0 100 100 none, DISTINCT_SCAN not reported
DocumentDB 0.114 gateway (MongoDB API) 50,000 50,000 hash aggregate (first group at ~349 ms) not exposed
DocumentDB 0.114 native (SQL function) 50,000 50,000 external merge sort, 14,144 kB (hash aggregation off) 1,768 blocks read, 1,771 written
DocumentDB 0.116 gateway (MongoDB API) 50,000 50,000 no, sorted streaming (first group at ~3 ms) not exposed
DocumentDB 0.116 native (SQL function) 50,000 50,000 none none
DocumentDB 0.116 native (SQL function), enableDistinctScanForGroupFirst on (default in 0.117) 100 100 none none

By default, 0.114 uses hash aggregation and blocks until it reads all 50,000 rows (first group around 349 ms). When hash aggregation is turned off, it sorts and spills 14 MB. Version 0.116 streams sorted groups directly from the index (first group in about 3 ms), avoiding sorting and temporary I/O. With default 0.116 settings, DocumentDB does not use the distinct-scan access path for this pipeline, so it still reads all 50,000 entries, but grouping is no longer blocking. Enabling enableDistinctScanForGroupFirst (the default in 0.117) brings it to MongoDB's profile: 100 keys and 100 documents.

Note that the official description of this optimization is "pushing suffix sort keys into the accumulator in $sortGroup", but $sortGroup is not a user-facing MongoDB aggregation operator. It's an internal stage name in DocumentDB's codebase: when a pipeline has a $sort immediately followed by a $group, the two are fused into a single internal stage.

Top comments (0)