|
14 | 14 | # KIND, either express or implied. See the License for the |
15 | 15 | # specific language governing permissions and limitations |
16 | 16 | # under the License. |
17 | | -"""Benchmark a single-value metrics predicate through manifest planning. |
| 17 | +"""Benchmark a realistic 15-leaf metrics predicate through manifest planning. |
18 | 18 |
|
19 | 19 | Run with: |
20 | 20 | uv run pytest tests/benchmark/test_metrics_evaluator_benchmark.py -v -s -m benchmark |
|
31 | 31 |
|
32 | 32 | import pyiceberg.table as table_module |
33 | 33 | from pyiceberg.conversions import to_bytes |
34 | | -from pyiceberg.expressions import EqualTo |
| 34 | +from pyiceberg.expressions import And, BooleanExpression, EqualTo, GreaterThanOrEqual, LessThanOrEqual, Or |
35 | 35 | from pyiceberg.manifest import DataFile, FileFormat, ManifestContent, ManifestFile |
36 | 36 | from pyiceberg.schema import Schema |
37 | 37 | from pyiceberg.table import ManifestGroupPlanner, Table |
|
41 | 41 |
|
42 | 42 |
|
43 | 43 | def _data_file(file_number: int) -> DataFile: |
44 | | - value = file_number % 11 |
45 | | - value_bytes = to_bytes(LongType(), value) |
| 44 | + event_day = file_number % 11 |
| 45 | + region_id = file_number % 15 |
| 46 | + event_day_bytes = to_bytes(LongType(), event_day) |
| 47 | + region_id_bytes = to_bytes(LongType(), region_id) |
46 | 48 | return DataFile.from_args( |
47 | 49 | file_path=f"s3://bucket/data-{file_number}.parquet", |
48 | 50 | file_format=FileFormat.PARQUET, |
49 | 51 | partition=Record(), |
50 | 52 | record_count=100, |
51 | 53 | file_size_in_bytes=1, |
52 | | - value_counts={1: 100}, |
53 | | - null_value_counts={1: 0}, |
54 | | - lower_bounds={1: value_bytes}, |
55 | | - upper_bounds={1: value_bytes}, |
| 54 | + value_counts={1: 100, 2: 100}, |
| 55 | + null_value_counts={1: 0, 2: 0}, |
| 56 | + lower_bounds={1: event_day_bytes, 2: region_id_bytes}, |
| 57 | + upper_bounds={1: event_day_bytes, 2: region_id_bytes}, |
56 | 58 | ) |
57 | 59 |
|
58 | 60 |
|
| 61 | +def _metrics_filter() -> BooleanExpression: |
| 62 | + """Select five day ranges, each scoped to a region.""" |
| 63 | + windows = ((0, 1, 1), (2, 3, 4), (4, 5, 7), (6, 7, 10), (8, 10, 13)) |
| 64 | + branches = [ |
| 65 | + And( |
| 66 | + And(GreaterThanOrEqual("event_day", start_day), LessThanOrEqual("event_day", end_day)), |
| 67 | + EqualTo("region_id", region_id), |
| 68 | + ) |
| 69 | + for start_day, end_day, region_id in windows |
| 70 | + ] |
| 71 | + |
| 72 | + combined = branches[0] |
| 73 | + for branch in branches[1:]: |
| 74 | + combined = Or(combined, branch) |
| 75 | + return combined |
| 76 | + |
| 77 | + |
59 | 78 | def _manifest_file(manifest_number: int) -> ManifestFile: |
60 | 79 | return ManifestFile.from_args( |
61 | 80 | manifest_path=f"s3://bucket/manifest-{manifest_number}.avro", |
@@ -87,12 +106,13 @@ def test_metrics_evaluator_reuse( |
87 | 106 | ) -> None: |
88 | 107 | num_files = 1_000 |
89 | 108 | schema = Schema( |
90 | | - NestedField(1, "x", LongType(), required=True), |
91 | | - *(NestedField(field_id, f"unused_{field_id}", LongType(), required=False) for field_id in range(2, 102)), |
| 109 | + NestedField(1, "event_day", LongType(), required=True), |
| 110 | + NestedField(2, "region_id", LongType(), required=True), |
| 111 | + *(NestedField(field_id, f"unused_{field_id}", LongType(), required=False) for field_id in range(3, 103)), |
92 | 112 | schema_id=table_v2.metadata.current_schema_id, |
93 | 113 | ) |
94 | 114 | metadata = table_v2.metadata.model_copy(update={"schemas": [schema]}) |
95 | | - planner = ManifestGroupPlanner(table_metadata=metadata, io=table_v2.io, row_filter=EqualTo("x", 5)) |
| 115 | + planner = ManifestGroupPlanner(table_metadata=metadata, io=table_v2.io, row_filter=_metrics_filter()) |
96 | 116 | data_files = [_data_file(file_number) for file_number in range(num_files)] |
97 | 117 | manifests = [_manifest_file(manifest_number) for manifest_number in range(0, num_files, files_per_manifest)] |
98 | 118 | files_by_manifest = { |
@@ -120,14 +140,14 @@ def open_manifest( |
120 | 140 | def evaluate_files() -> int: |
121 | 141 | return sum(len(entries) for entries in planner.plan_manifest_entries(manifests)) |
122 | 142 |
|
123 | | - assert evaluate_files() == (91 if partition_matches else 0) |
| 143 | + assert evaluate_files() == (67 if partition_matches else 0) |
124 | 144 | number = 1 if partition_matches else (1_000 if files_per_manifest == 1_000 else 10) |
125 | 145 | timings_ms = [timing * 1_000 / number for timing in timeit.repeat(evaluate_files, number=number, repeat=5)] |
126 | 146 | file_label = "file" if files_per_manifest == 1 else "files" |
127 | 147 | pruning_label = "all files pruned" if not partition_matches else "metrics evaluated" |
128 | 148 |
|
129 | 149 | print( |
130 | 150 | f"Evaluated metrics for {num_files} files with {files_per_manifest} {file_label} per manifest, " |
131 | | - f"a 102-column schema, and a single-value predicate ({pruning_label}) in " |
| 151 | + f"a 102-column schema, and a 15-leaf predicate ({pruning_label}) in " |
132 | 152 | f"{statistics.mean(timings_ms):.3f}ms (best: {min(timings_ms):.3f}ms)" |
133 | 153 | ) |
0 commit comments