Skip to content

Add batched table stats collection spark app - #789

Merged
abhisheknath2011 merged 4 commits into
linkedin:mainfrom
abhisheknath2011:batched-sc
Oct 6, 2026
Merged

abhisheknath2011 merged 4 commits into
linkedin:mainfrom
abhisheknath2011:batched-sc

Conversation

@abhisheknath2011

Copy link
Copy Markdown
Member

Summary

Add BatchedTableStatsCollectionSparkApp — the multi-table Spark app that runs table stats collection for a bin of tables in a single job, mirroring BatchedOrphanFilesDeletionSparkApp. This is the execution side for the optimizer's stats-collection operation: the scheduler bin-packs TABLE_STATS_COLLECTION ops (by file count) and submits one job per bin; this app processes that bin.

Motivation

The optimizer can now analyze and bin-pack TABLE_STATS_COLLECTION operations, but there was no batched Spark app to execute a bin. The single-table TableStatsCollectionSparkApp remains the unit when bin size is 1; this app processes many (table, operationId) pairs in one job so bin-packing actually reduces the number of Spark jobs.

Changes

BatchedTableStatsCollectionSparkApp (new)

  • Extends BaseSparkApp; one job processes a list of (fqtn, operationId, tableUuid) the scheduler packed into a bin.
  • CLI: --tableNames, --operationIds, --tableUuids (parallel CSV lists), --resultsEndpoint, --driverParallelism.
  • A fixed thread pool runs one worker per table; each worker collects and publishes the same four artifacts as the single-table app (collectTableStats + commit events + partition events + partition stats), then posts the per-operation outcome to the Optimizer Service with operationType = TABLE_STATS_COLLECTION.
  • Per-table failure isolation: a worker catches Throwable, reports FAILED for its own operation, and the job continues; the job exits 0 if ≥1 table succeeds and throws only if all fail.
  • Core-stats semantics: a null/throwing collectTableStats marks that table FAILED (so the analyzer's failure cadence retries it); commit-/partition-level artifacts are best-effort.
  • Defensive result reporting mirrors the OFD app: a missed completion callback leaves the row at SCHEDULED (logged + counted) for the analyzer's stale-timeout rather than silently dropping.

Supporting changes

  • api/spec/OperationType: add TABLE_STATS_COLLECTION, so the generated optimizer client's
    UpdateOperationRequest.OperationTypeEnum includes it for reportResult.
  • AppConstants: add STATS_MAX_BATCH_SIZE = 100 — a footgun guard against an oversized batch OOMing the driver (the scheduler's per-bin table cap defaults to 25; this is the hard ceiling).

Intentional differences from BatchedOrphanFilesDeletionSparkApp

  • No post-job table-state validation — stats collection does not mutate the table.
  • No OFD-specific options (ttl / backupDir / concurrentDeletes / streamResults).

Issue] Briefly discuss the summary of the changes made in this
pull request in 2-3 lines.

  • Client-facing API Changes
  • Internal API Changes
  • Bug Fixes
  • New Features
  • Performance Improvements
  • Code Style
  • Refactoring
  • Documentation
  • Tests

For all the boxes checked, please include additional details of the changes made in this pull request.

Testing Done

./gradlew clean build passed

  • BatchedTableStatsCollectionSparkAppArgsTest (pure-Java, 7 cases): parallel-list parsing,
    whitespace trimming, optional lists, mismatched-length rejection, non-qualified name rejection,
    and the STATS_MAX_BATCH_SIZE guard.
  • :apps:openhouse-spark-apps_2.12 compiles (the optimizer client regenerates from the updated spec
    with TABLE_STATS_COLLECTION) and the test suite passes (Java 17).
  • Manually Tested on local docker setup. Please include commands ran, and their output.
  • Added new tests for the changes made.
  • Updated existing tests to reflect the changes made.
  • No tests added or updated. Please explain why. If unsure, please feel free to ask for help.
  • Some other form of testing like staging or soak time in production. Please explain.

For all the boxes checked, include a detailed description of the testing done for the changes made in this pull request.

Additional Information

  • Breaking Changes
  • Deprecations
  • Large PR broken into smaller PRs, and PR plan linked in the description.

For all the boxes checked, include additional details of the changes made in this pull request.

@mkuchenbecker mkuchenbecker left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What are your thoughts on the app having a metrics sink rather than extending and overriding a function that logs? Then we can test it end to end with a generic sink and this app is functionally complete rather than a logger.

@abhisheknath2011

abhisheknath2011 commented Oct 6, 2026 •

Copy link
Copy Markdown
Member Author

What are your thoughts on the app having a metrics sink rather than extending and overriding a function that logs? Then we can test it end to end with a generic sink and this app is functionally complete rather than a logger.

Good call - switched from overridable logging methods to an injected sink.

  • Added a  StatsCollectionSink  interface ( publishStats  /  publishCommitEvents  /  publishPartitionEvents  /  publishPartitionStats ); the app now has a sink and publishes to it instead of overriding  publish*  methods.
  • Default  LoggingStatsCollectionSink  keeps it runnable out of the box, so the app is functionally complete; a deployment injects a durable sink (e.g. Kafka) via composition, no subclassing.
  • The ITest now drives the real app end-to-end through a generic capturing sink and asserts what was published.

@abhisheknath2011
abhisheknath2011 merged commit edcae62 into linkedin:main Oct 6, 2026
1 check passed
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.

2 participants