Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,9 @@ jobs:
- name: Check Iceberg shard inventory validation
run: python3 dev/ci/test-iceberg-shards.py

- name: Check Iceberg write report summary
run: python3 dev/ci/test-summarize-iceberg-writes.py

- name: Check CI config invariants
run: python3 dev/ci/check-ci-config.py

Expand Down
32 changes: 30 additions & 2 deletions .github/workflows/iceberg_spark_test_reusable.yml
Original file line number Diff line number Diff line change
Expand Up @@ -154,12 +154,20 @@ jobs:
run: |
cd apache-iceberg
rm -rf /root/.m2/repository/org/apache/parquet # somehow parquet cache requires cleanups
ENABLE_COMET=true ENABLE_COMET_ONHEAP=true ./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \
# COMET_ICEBERG_WRITE_REPORT_DIR records which writer ran each Iceberg
# write; see dev/ci/summarize-iceberg-writes.py.
ENABLE_COMET=true ENABLE_COMET_ONHEAP=true COMET_ICEBERG_WRITE_REPORT_DIR="$PWD/build/comet-iceberg-writes" \
./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \
:iceberg-spark:iceberg-spark-${{ inputs.spark-short }}_${{ inputs.scala }}:test \
--init-script ../dev/ci/iceberg-test-shards.gradle \
-PcometShardTask=:iceberg-spark:iceberg-spark-${{ inputs.spark-short }}_${{ inputs.scala }}:test \
-PcometShardIndex=${{ matrix.shard }} -PcometShardCount=${{ needs.build-native.outputs.shard-count }} \
-Pquick=true -x javadoc
- name: Summarize Iceberg writes
if: ${{ !cancelled() }}
run: |
python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark shard ${{ matrix.shard }}" \
apache-iceberg/build/comet-iceberg-writes
- name: Upload Iceberg shard inventory and test reports
if: ${{ !cancelled() }}
# iceberg-spark-shard-coverage downloads the inventory, so a flaky
Expand All @@ -170,6 +178,7 @@ jobs:
path: |
apache-iceberg/**/build/comet-shards/*.json
apache-iceberg/**/build/test-results/test/*.xml
apache-iceberg/build/comet-iceberg-writes/*.jsonl
retention-days: 7

iceberg-spark-shard-coverage:
Expand All @@ -190,6 +199,11 @@ jobs:
run: |
python3 dev/ci/check-iceberg-shards.py --manifests iceberg-shard-reports \
--task :iceberg-spark:iceberg-spark-${{ inputs.spark-short }}_${{ inputs.scala }}:test
- name: Summarize Iceberg writes across shards
if: ${{ !cancelled() }}
run: |
python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark, all shards" \
iceberg-shard-reports

iceberg-spark-extensions:
needs: build-native
Expand Down Expand Up @@ -222,9 +236,23 @@ jobs:
run: |
cd apache-iceberg
rm -rf /root/.m2/repository/org/apache/parquet # somehow parquet cache requires cleanups
ENABLE_COMET=true ENABLE_COMET_ONHEAP=true ./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \
ENABLE_COMET=true ENABLE_COMET_ONHEAP=true COMET_ICEBERG_WRITE_REPORT_DIR="$PWD/build/comet-iceberg-writes" \
./gradlew -DsparkVersions=${{ inputs.spark-short }} -DscalaVersion=${{ inputs.scala }} -DflinkVersions= -DkafkaVersions= \
:iceberg-spark:iceberg-spark-extensions-${{ inputs.spark-short }}_${{ inputs.scala }}:test \
-Pquick=true -x javadoc
- name: Summarize Iceberg writes
if: ${{ !cancelled() }}
run: |
python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark-extensions" \
apache-iceberg/build/comet-iceberg-writes
- name: Upload Iceberg write report
if: ${{ !cancelled() }}
uses: ./.github/actions/upload-artifact-retry
with:
name: iceberg-spark-extensions-writes-${{ inputs.iceberg-full }}-spark-${{ inputs.spark-full }}-scala-${{ inputs.scala }}-jdk${{ inputs.java }}-attempt-${{ github.run_attempt }}
path: apache-iceberg/build/comet-iceberg-writes/*.jsonl
if-no-files-found: ignore
retention-days: 7

iceberg-spark-runtime:
needs: build-native
Expand Down
171 changes: 171 additions & 0 deletions dev/ci/summarize-iceberg-writes.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,171 @@
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.

"""Summarize which writer ran the Iceberg writes of an Iceberg Spark test run.

The Iceberg Spark test jobs set COMET_ICEBERG_WRITE_REPORT_DIR, so Comet's
IcebergWriteReportListener appends one JSON line per Iceberg write to a file
in that directory. This prints how many writes ran on Comet's native writer,
how many Comet's split operator left on Iceberg's JVM writer and why, and how
many Spark planned without Comet's split operator:

python3 dev/ci/summarize-iceberg-writes.py --title "iceberg-spark shard 1" DIR...

Files under a directory named like a shard artifact (...-shard-N-attempt-M)
are counted only for the latest attempt of each shard, so a rerun of failed
jobs does not count a shard twice. The latest attempt is the newest such
directory, whether or not it holds any report files, so a rerun that recorded
no writes is reported as missing rather than replaced by an earlier attempt.
The summary is also appended to $GITHUB_STEP_SUMMARY when that is set. It
never fails the job.
"""

import argparse
from collections import Counter
import json
import os
from pathlib import Path
import re


SHARD_ATTEMPT = re.compile(r"-shard-(\d+)-attempt-(\d+)$")
WRITERS = [
("native", "Comet native writer"),
("jvm", "Iceberg JVM writer under Comet's split operator"),
("spark", "Spark V2 write, not planned by Comet's split operator"),
]
TOP_REASONS = 20


def shard_attempt(path):
"""The (shard, attempt) of the shard artifact directory holding path, or None."""
for part in path.parts:
match = SHARD_ATTEMPT.search(part)
if match:
return int(match.group(1)), int(match.group(2))
return None


def report_files(roots):
"""Every report file under roots, keeping only the latest attempt of each shard.

Returns the files and the (shard, attempt) pairs whose latest attempt holds no report file.
The latest attempt comes from the artifact directories rather than the report files, because
an attempt that recorded no writes still uploads its shard inventory and test reports.
"""
latest = {}
files = []
for root in map(Path, roots):
for path in [root, *sorted(root.rglob("*"))]:
key = shard_attempt(path)
if key:
latest[key[0]] = max(latest.get(key[0], 0), key[1])
if path.suffix == ".jsonl" and path.is_file():
files.append((path, key))
kept = [(path, key) for path, key in files if key is None or latest[key[0]] == key[1]]
reported = {key for _, key in kept}
missing = sorted(key for key in latest.items() if key not in reported)
return [path for path, _ in kept], missing


def load(roots):
files, missing = report_files(roots)
writes = []
for path in files:
for line in path.read_text(encoding="utf-8").splitlines():
if line.strip():
writes.append(json.loads(line))
return writes, missing


def cell(text):
return " ".join(text.split()).replace("|", "\\|")


def summarize(title, writes, missing=()):
lines = [f"### Iceberg writes: {title}", ""]
for shard, attempt in missing:
lines += [
f"Shard {shard} recorded no Iceberg writes in its latest attempt ({attempt}), "
"so none of its writes are counted below.",
"",
]
if not writes:
lines.append(
"No Iceberg writes were recorded. Either the target ran none or "
"COMET_ICEBERG_WRITE_REPORT_DIR did not reach the test JVMs."
)
return "\n".join(lines) + "\n"

total = len(writes)
by_writer = Counter(w["writer"] for w in writes)
lines += ["| Writer | Writes | Share |", "| --- | ---: | ---: |"]
for key, label in WRITERS:
count = by_writer.get(key, 0)
lines.append(f"| {label} | {count} | {100.0 * count / total:.1f}% |")
lines += [f"| Total | {total} | |", ""]
failed = sum(1 for w in writes if w.get("failed"))
if failed:
lines += [f"{failed} of the {total} writes ran in queries that failed.", ""]

reasons = Counter()
for w in writes:
if w["writer"] == "jvm":
for reason in w.get("reasons") or ["(no reason recorded)"]:
reasons[reason] += 1
if reasons:
lines += [
"#### Why the split operator kept the JVM writer",
"",
"A write can have several reasons, so the counts can add up to more than the "
"JVM writes.",
"",
"| Writes | Reason |",
"| ---: | --- |",
]
for reason, count in reasons.most_common(TOP_REASONS):
lines.append(f"| {count} | {cell(reason)} |")
if len(reasons) > TOP_REASONS:
lines.append(f"| | and {len(reasons) - TOP_REASONS} more reasons |")
lines.append("")

operators = Counter(w["node"] for w in writes if w["writer"] == "spark")
if operators:
lines += ["#### Spark V2 writes by operator", "", "| Writes | Operator |", "| ---: | --- |"]
for node, count in operators.most_common():
lines.append(f"| {count} | {cell(node)} |")
lines.append("")

return "\n".join(lines) + "\n"


def main():
parser = argparse.ArgumentParser(description=__doc__.splitlines()[0])
parser.add_argument("--title", required=True, help="heading for the summary")
parser.add_argument("roots", nargs="+", help="directories holding the report files")
args = parser.parse_args()

summary = summarize(args.title, *load(args.roots))
print(summary)
step_summary = os.environ.get("GITHUB_STEP_SUMMARY")
if step_summary:
with open(step_summary, "a", encoding="utf-8") as out:
out.write(summary + "\n")


if __name__ == "__main__":
main()
105 changes: 105 additions & 0 deletions dev/ci/test-summarize-iceberg-writes.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
#!/usr/bin/env python3
#
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.

"""Fast regression tests for the Iceberg write report summary."""

import importlib.util
import json
from pathlib import Path
import tempfile
import unittest


SPEC = importlib.util.spec_from_file_location(
"summarize_iceberg_writes", Path(__file__).with_name("summarize-iceberg-writes.py"))
SUMMARIZE = importlib.util.module_from_spec(SPEC)
SPEC.loader.exec_module(SUMMARIZE)

ARTIFACT = "iceberg-spark-1.11.0-spark-4.1.3-scala-2.13-jdk17-shard-{}-attempt-{}"


def records(*writers):
return "".join(
json.dumps({"writer": w, "node": "AppendData", "reasons": [], "failed": False}) + "\n"
for w in writers)


class SummarizeIcebergWritesTest(unittest.TestCase):
def setUp(self):
self.temp = tempfile.TemporaryDirectory(prefix="comet-iceberg-writes-test-")
self.addCleanup(self.temp.cleanup)
self.root = Path(self.temp.name)

def attempt(self, shard, attempt, *writers):
"""A shard attempt's artifact as the coverage job downloads it.

Every attempt uploads its test reports. Only an attempt that recorded writes has a
report file.
"""
artifact = self.root / ARTIFACT.format(shard, attempt)
reports = artifact / "build/test-results/test"
reports.mkdir(parents=True)
(reports / "TEST-org.example.TestFixture.xml").write_text("<testsuite/>")
if writers:
writes = artifact / "build/comet-iceberg-writes"
writes.mkdir(parents=True)
(writes / "iceberg-writes-fixture.jsonl").write_text(records(*writers))

def load(self):
writes, missing = SUMMARIZE.load([self.root])
return sorted(w["writer"] for w in writes), missing

def test_every_shard_counts_once(self):
self.attempt(1, 1, "native")
self.attempt(2, 1, "jvm", "spark")
self.assertEqual(self.load(), (["jvm", "native", "spark"], []))

def test_latest_attempt_replaces_an_earlier_one(self):
self.attempt(1, 1, "native", "native")
self.attempt(1, 2, "jvm")
self.assertEqual(self.load(), (["jvm"], []))

def test_retry_without_writes_is_reported_missing_instead_of_stale(self):
self.attempt(1, 1, "native")
self.attempt(1, 2)
self.attempt(2, 1, "spark")
self.assertEqual(self.load(), (["spark"], [(1, 2)]))
summary = SUMMARIZE.summarize("fixture", *SUMMARIZE.load([self.root]))
self.assertIn("Shard 1 recorded no Iceberg writes in its latest attempt (2)", summary)
self.assertIn("| Comet native writer | 0 | 0.0% |", summary)
self.assertIn("| Spark V2 write, not planned by Comet's split operator | 1 | 100.0% |",
summary)

def test_root_that_is_itself_a_shard_artifact(self):
self.attempt(1, 3)
writes, missing = SUMMARIZE.load([self.root / ARTIFACT.format(1, 3)])
self.assertEqual((writes, missing), ([], [(1, 3)]))

def test_files_outside_shard_artifacts_always_count(self):
# A shard job and dev/local-ci.sh summarize their own report directory directly.
(self.root / "iceberg-writes-local.jsonl").write_text(records("native", "jvm"))
self.assertEqual(self.load(), (["jvm", "native"], []))

def test_no_writes_at_all(self):
summary = SUMMARIZE.summarize("fixture", *SUMMARIZE.load([self.root / "missing"]))
self.assertIn("No Iceberg writes were recorded", summary)


if __name__ == "__main__":
unittest.main()
11 changes: 10 additions & 1 deletion dev/local-ci.sh
Original file line number Diff line number Diff line change
Expand Up @@ -389,16 +389,25 @@ run_iceberg() {
gradlew ":iceberg-spark:iceberg-spark-runtime-${SPARK}_${SCALA}:integrationTest"
;;
esac
case "$target" in
shard-* | extensions)
python3 "$REPO/dev/ci/summarize-iceberg-writes.py" --title "$target" \
"$dest/build/comet-iceberg-writes/$target"
;;
esac
ok "$target took $(hms $((SECONDS - started)))"
done
}

# Reads $dest, $spark and $SCALA from run_iceberg.
# Reads $dest, $spark, $SCALA and $target from run_iceberg.
gradlew() {
(
cd "$dest"
# shellcheck disable=SC2031
export SPARK_LOCAL_IP=localhost ENABLE_COMET=true ENABLE_COMET_ONHEAP=true
# One directory per target, emptied first, so a rerun reports only its own writes.
export COMET_ICEBERG_WRITE_REPORT_DIR="$dest/build/comet-iceberg-writes/$target"
rm -rf "$COMET_ICEBERG_WRITE_REPORT_DIR"
./gradlew "-DsparkVersions=$SPARK" "-DscalaVersion=$SCALA" \
-DflinkVersions= -DkafkaVersions= "$@" -Pquick=true -x javadoc
)
Expand Down
Loading
Loading