{"record":{"id":"5330cf297632d697","repo":"apache/iceberg","slug":"file-at-offset-contains-records-exceedin-5330cf","errorCode":null,"errorMessage":"File {} at offset {} contains {} records, exceeding maxRecordsPerMicroBatch limit of {}. This file will be processed entirely to guarantee forward progress. Consider increasing the limit or writing smaller files to avoid unexpected memory usage.","messagePattern":"File (.+?) at offset (.+?) contains (.+?) records, exceeding maxRecordsPerMicroBatch limit of (.+?)\\. This file will be processed entirely to guarantee forward progress\\. Consider increasing the limit or writing smaller files to avoid unexpected memory usage\\.","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/AsyncSparkMicroBatchPlanner.java","lineNumber":328,"sourceCode":"        if (filesSeen == 0) {\n          return null;\n        }\n        LOG.debug(\n            \"latestOffset hit file limit at {}, rows: {}, files: {}\",\n            elem.first(),\n            rowsSeen,\n            filesSeen);\n        return elem.first();\n      }\n\n      // Soft limit on rows - include file FIRST, then check\n      rowsSeen += fileRows;\n      filesSeen += 1;\n\n      // Check if we've hit the row limit after including this file\n      if (rowsSeen >= unpackedLimits.getMaxRows()) {\n        if (filesSeen == 1 && rowsSeen > unpackedLimits.getMaxRows()) {\n          LOG.warn(\n              \"File {} at offset {} contains {} records, exceeding maxRecordsPerMicroBatch limit of {}. \"\n                  + \"This file will be processed entirely to guarantee forward progress. \"\n                  + \"Consider increasing the limit or writing smaller files to avoid unexpected memory usage.\",\n              elem.second().file().location(),\n              elem.first(),\n              fileRows,\n              unpackedLimits.getMaxRows());\n        }\n        // Return the offset of the NEXT element (or synthesize tail+1)\n        if (i + 1 < queueSnapshot.size()) {\n          LOG.debug(\n              \"latestOffset hit row limit at {}, rows: {}, files: {}\",\n              queueSnapshot.get(i + 1).first(),\n              rowsSeen,\n              filesSeen);\n          return queueSnapshot.get(i + 1).first();\n        } else {\n          // This is the last element - return tail+1","sourceCodeStart":310,"sourceCodeEnd":346,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/source/AsyncSparkMicroBatchPlanner.java#L310-L346","documentation":"AsyncSparkMicroBatchPlanner limits each micro-batch by maxRecordsPerMicroBatch. If a single file alone exceeds the row limit, including it wholly is the only way to guarantee forward progress; the planner logs this warning and processes the oversized file in one batch, potentially raising memory usage.","triggerScenarios":"computeLimitedOffset (called from latestOffset) encounters a first file (filesSeen == 1) whose row count exceeds unpackedLimits.getMaxRows() from the maxRecordsPerMicroBatch limit.","commonSituations":"Downstream writers producing very large files relative to a small maxRecordsPerMicroBatch setting; users lowering the limit for latency without considering file sizes.","solutions":["Increase maxRecordsPerMicroBatch so it comfortably exceeds the largest expected file's row count","Configure upstream writers to write smaller files (e.g. smaller target file size or more frequent compaction)","Compact/rewrite large files so future batches can respect the limit","Monitor memory during such batches and size executors accordingly"],"exampleFix":"// before\n.option(\"maxRecordsPerMicroBatch\", \"1000\")   // files contain >1000 rows each\n// after\n.option(\"maxRecordsPerMicroBatch\", \"500000\") // above largest file size","handlingStrategy":"validation","validationCode":"long maxRows = Long.parseLong(limits.get(\"maxRecordsPerMicroBatch\"));\nlong largestFileRows = files.stream().mapToLong(f -> f.recordCount()).max().orElse(0);\nif (largestFileRows > maxRows) { /* raise maxRecordsPerMicroBatch above largestFileRows */ }","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Keep maxRecordsPerMicroBatch well above the largest data file's record count","Configure writers with smaller target file sizes","Compact oversized files proactively","Alert on file record counts relative to streaming batch limits"],"tags":["spark","streaming","micro-batch","memory"],"backgroundTag":"value-out-of-range","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}