Skip to content

fix: stream mergeResultFilesToParquet one file at a time - #117

Merged
abhinav-pola merged 3 commits into
mainfrom
devin/1790357465-streaming-merge
Sep 25, 2026
Merged

abhinav-pola merged 3 commits into
mainfrom
devin/1790357465-streaming-merge

Conversation

@abhinav-pola

@abhinav-pola abhinav-pola commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

TL;DR

mergeResultFilesToParquet now holds one decoded input file at a time and writes to a caller-supplied parquet Writer, so merging a large run no longer needs every row plus the whole output buffer in memory.

What changed?

  • Breaking signature change. The only consumer is services/gcp-harbor-run/src/publish-native-run.ts in openrouter-web, which moves over in the subtree-pull PR.
// before
mergeResultFilesToParquet(files: readonly (readonly BenchmarkResultRow[])[], meta): Buffer
// after
type ResultFileReader = () => Promise<readonly BenchmarkResultRow[]>;
mergeResultFilesToParquet({ files: readonly ResultFileReader[], meta, writer: Writer }):
  Promise<{ summary: ChunkResultSummary | null; benchmarkConfig: string | null }>
  • It makes two passes over the readers, one file at a time:
    1. Pass 1 collects row counts, light SampleScores (id, epoch, value), summed usage, the weighted primary score, and the first file's epochs, temperature and benchmark_config.
    2. Pass 2 re-reads each non-empty file and checks it against pass 1: row count, each row's sample_id/epoch/score_value, and a fingerprint of the first row's run-level fields. Any mismatch throws. It then overwrites the run-level columns and writes rows through ParquetWriter in row groups of at most 100 rows.
  • The merge semantics and the full raw schema stay the same. Empty files are still skipped, extra_scores is still null, and format_version is still RESULT_FORMAT_VERSION. The schema and KV metadata match runResultToParquet.
  • The returned summary equals summarizeChunkRows(mergedRows), so callers don't need to decode the output again. summarizeChunkRows and the merge now share summarizeSampleScores.
  • rowsToSampleScores no longer copies answer, because it only feeds aggregates.
  • New export ./parquet-stream-writer. parquetStreamWriter(writeChunk, chunkBytes = 1_000_000) is a Writer that hands bytes to writeChunk in order on each row-group flush and on finish, instead of retaining the file. It's meant for piping into a GCS upload stream.

Why?

Providers need one GCS download link per run. The only existing merge built everything in memory: all decoded rows, the full output Buffer, and then the caller decoded it again. That will OOM on agentic runs with hundreds of trials.

How to test

Memory probe, run locally: 60 files × 20 rows, each row with a unique 200 KB explanation, 222 MB of parquet on disk. The old path mirrors publish-native-run (read all, merge to a Buffer, parse again). The new path uses file readers plus parquetStreamWriter appending to disk.

old (origin/main):  max RSS 2,445,668 KB   1200 rows   output 231,805,554 B
new (this branch):  max RSS   285,208 KB   1200 rows   output 232,037,709 B

The new output is about 0.1% larger because each row group carries its own dictionary and page index. These numbers are from the first commit (one row group per file). The unit tests check that the streamed bytes match an in-memory ByteWriter merge byte-for-byte, and that readers are called sequentially, twice each.

Reviewer focus

  • Row groups never span input files and hold at most 100 rows. Many tiny trial files mean many small row groups, and footer metadata grows linearly with the row-group count.
  • parquetStreamWriter overrides ensure, flush and finish on a ByteWriter, the same pattern as hyparquet-writer's own fileWriter. getBuffer and getBytes throw on purpose.

Checklist

  • Tests cover changed behavior
  • Public API or configuration changes are backward compatible, or the break is documented
  • Benchmark changes document dataset provenance and licensing (n/a)
  • No credentials, private results, or restricted dataset contents are included
  • Documentation is updated where needed (n/a)

Link to Devin session: https://openrouter.devinenterprise.com/sessions/cff9171c8dd8432ebb515e2a54374727
Open in Devin Desktop: https://openrouter.devinenterprise.com/desktop/session/cff9171c8dd8432ebb515e2a54374727?variant=devin
Requested by: @abhinav-pola


Devin Review

@devin-ai-integration

Copy link
Copy Markdown
Contributor

I'll fix CI failures and address comments from users with write access that start with 'DevinAI' or '@devin'.

  • Disable automatic comment, CI, and merge conflict monitoring

Original prompt from Abhinav

SYSTEM:
<latest_message>
Abhinav Pola (U090K0G7JF3) [ts=1790356866.277609]: @Devin read this thread <https://openrouter.slack.com/archives/C0BMHG5CG1E/p1789516956636439>
</latest_message>

=== BEGIN THREAD HISTORY (in #agents-benchmarks) ===
Abhinav Pola (U090K0G7JF3) [ts=1790356866.277609]: @Devin read this thread <https://openrouter.slack.com/archives/C0BMHG5CG1E/p1789516956636439>
=== END THREAD HISTORY ===
Channel ID: C0BAKP8P5C3
Thread URL: https://openrouter.slack.com/archives/C0BAKP8P5C3/p1790356866277609?thread_ts=1790356866.277609&amp;cid=C0BAKP8P5C3

The <latest_message> is the message that you should use to guide your goals + task for this session, and you should use the rest of the slack thread as context.
A [ts=...] marker on a Slack message is that message's timestamp. To act on a specific message with the slack tool (e.g. adding an emoji reaction via the reaction command), pass that value as timestamp along with the Channel ID — no extra lookup call is needed.

devin-ai-integration[bot]

This comment was marked as resolved.

@abhinav-pola
abhinav-pola merged commit 801508d into main Sep 25, 2026
5 checks passed
@abhinav-pola
abhinav-pola deleted the devin/1790357465-streaming-merge branch September 25, 2026 17:48
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant