diff --git a/AGENTS.md b/AGENTS.md index f0bf34ea..18a15049 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -26,6 +26,8 @@ Instructions for coding agents working in IceGraph. unless the user explicitly requests it. - Do not disable lint, type, or formatting rules to make checks pass. - Python files must not exceed 400 lines. This is a repository convention, not a formatter check. +- In Python, use an explicit `for` loop with `break` for first-match searches. Do not use + `next(...)` with a generator expression for these lookups. - Define every backend runtime environment setting in `backend/env.py`. - Never write docstrings, or comments in Python code, unless the user asks for them. Exception: follow an existing per-entry comment convention, such as the one-line comment above each setting in `backend/env.py`. diff --git a/backend/collectors/collect_partition_statistics.py b/backend/collectors/collect_partition_statistics.py index 0a12931d..c6f53079 100644 --- a/backend/collectors/collect_partition_statistics.py +++ b/backend/collectors/collect_partition_statistics.py @@ -35,30 +35,39 @@ def _collect_statistics_files(self, metadata_file_to_added_entries: Dict[str, Li statistics_files = self._build_statistics_files(metadata_file_to_added_entries) self._collect_partition_summaries(statistics_files) + def collect_file(self, metadata_path: str, entry: dict, include_samples: bool = True) -> PartitionStatisticsFileRecord: + statistics_file = self._parse_statistics_entry(metadata_path, entry) + self._collect_partition_summaries({statistics_file.file_path: statistics_file}, include_samples) + return statistics_file + def _build_statistics_files(self, metadata_file_to_added_entries: Dict[str, List[dict]]) -> Dict[str, PartitionStatisticsFileRecord]: statistics_files = {} for metadata_file in self._metadata_files: for entry in metadata_file_to_added_entries.get(metadata_file.file_path, []): - statistics_file = PartitionStatisticsFileRecord( - type=FileType.PARTITION_STATISTICS, - file_path=entry["statistics-path"], - child_files=[], - snapshot_id=entry["snapshot-id"], - file_size_in_bytes=str(entry["file-size-in-bytes"]), - partitions_count=None, - partitions_with_deletes=None, - partition_distribution={}, - sampled_partitions=[], - hidden_statistics_data=HiddenStatisticsMetadata(added_by_metadata_file=metadata_file.file_path), - ) + statistics_file = self._parse_statistics_entry(metadata_file.file_path, entry) statistics_files[statistics_file.file_path] = statistics_file self._statistics_files.append(statistics_file) return statistics_files - def _collect_partition_summaries(self, statistics_files: Dict[str, PartitionStatisticsFileRecord]) -> None: + @staticmethod + def _parse_statistics_entry(metadata_path: str, entry: dict) -> PartitionStatisticsFileRecord: + return PartitionStatisticsFileRecord( + type=FileType.PARTITION_STATISTICS, + file_path=entry["statistics-path"], + child_files=[], + snapshot_id=entry["snapshot-id"], + file_size_in_bytes=str(entry["file-size-in-bytes"]), + partitions_count=None, + partitions_with_deletes=None, + partition_distribution={}, + sampled_partitions=[], + hidden_statistics_data=HiddenStatisticsMetadata(added_by_metadata_file=metadata_path), + ) + + def _collect_partition_summaries(self, statistics_files: Dict[str, PartitionStatisticsFileRecord], include_samples: bool = True) -> None: try: - rows = PartitionStatisticsExtractor(self._table_name, list(statistics_files.values())).extract_dataframe().collect() + rows = PartitionStatisticsExtractor(self._table_name, list(statistics_files.values())).extract_dataframe(include_samples).collect() except Exception as e: logger.error(f"[{self._table_name}] Partition statistics batch read error", exc_info=True) for statistics_file in statistics_files.values(): @@ -76,7 +85,7 @@ def _collect_partition_summaries(self, statistics_files: Dict[str, PartitionStat metric: value for metric, value in (statistics_row["summary"]["partition_distribution"] or {}).items() if value is not None } sampled_partitions = sorted( - statistics_row["sampled_partitions"], + statistics_row.get("sampled_partitions", []), key=lambda partition: partition["last_updated_at"] if partition.get("last_updated_at") is not None else float("-inf"), reverse=True, ) diff --git a/backend/collectors/collect_table_metadata.py b/backend/collectors/collect_table_metadata.py index f17c9b14..f415097f 100644 --- a/backend/collectors/collect_table_metadata.py +++ b/backend/collectors/collect_table_metadata.py @@ -7,6 +7,8 @@ from pyspark.sql.types import ArrayType, MapType, StringType, StructType from base_classes.utils import collect_graph_metadata_file, format_snapshot_summary, timed +from collectors.collect_partition_statistics import CollectPartitionStatistics, PartitionStatisticsFileRecord +from graph_normalizer.utils import to_json_safe from spark_connect import open_spark_connect_session @@ -20,7 +22,7 @@ def collect_latest(self) -> dict[str, Any]: return self.collect(metadata_path) - def collect(self, metadata_path: str) -> dict[str, Any]: + def collect(self, metadata_path: str, partition_statistics_files: list[PartitionStatisticsFileRecord] | None = None) -> dict[str, Any]: spark = open_spark_connect_session() metadata_df = spark.read.option("multiLine", True).json(metadata_path) row = ( @@ -28,7 +30,6 @@ def collect(self, metadata_path: str) -> dict[str, Any]: .drop("metadata-log") .drop("snapshot-log") .drop("snapshots") - .drop("statistics") .first() ) @@ -36,8 +37,56 @@ def collect(self, metadata_path: str) -> dict[str, Any]: metadata["schemas"] = metadata.get("schemas", []) self._parse_schema_field_types(metadata["schemas"]) self._format_current_snapshot_summary(metadata["current-snapshot"]) + self._collect_current_partition_statistics(metadata, metadata_path, partition_statistics_files or []) - return {"table-name": self._table_name, "metadata_file_path": metadata_path, **metadata} + return to_json_safe({"table-name": self._table_name, "metadata_file_path": metadata_path, **metadata}) + + def _collect_current_partition_statistics( + self, metadata: dict[str, Any], metadata_path: str, partition_statistics_files: list[PartitionStatisticsFileRecord] + ) -> None: + statistics_file = self._collect_current_partition_statistics_file(metadata, metadata_path, partition_statistics_files) + if statistics_file is None: + return + + statistics = statistics_file.to_dict() + metadata["current-partition-statistics"] = { + key: statistics[key] + for key in ( + "file_path", + "snapshot_id", + "file_size_in_bytes", + "partitions_count", + "partitions_with_deletes", + "partition_distribution", + "errors", + "warnings", + ) + } + + def _collect_current_partition_statistics_file( + self, metadata: dict[str, Any], metadata_path: str, partition_statistics_files: list[PartitionStatisticsFileRecord] + ) -> PartitionStatisticsFileRecord | None: + snapshot_id = metadata.get("current-snapshot-id") + if snapshot_id is None or snapshot_id == -1: + return None + + entry = None + for candidate in metadata.get("partition-statistics") or []: + if candidate["snapshot-id"] == snapshot_id: + entry = candidate + break + if entry is None: + return None + + statistics_file = None + for candidate in partition_statistics_files: + if candidate.file_path == entry["statistics-path"]: + statistics_file = candidate + break + if statistics_file is None: + statistics_file = CollectPartitionStatistics(self._table_name, []).collect_file(metadata_path, entry, include_samples=False) + + return statistics_file @staticmethod def _current_snapshot_column(metadata_df: DataFrame) -> Column: diff --git a/backend/extractors/partition_statistics_extractor.py b/backend/extractors/partition_statistics_extractor.py index 0901d4f5..495b6012 100644 --- a/backend/extractors/partition_statistics_extractor.py +++ b/backend/extractors/partition_statistics_extractor.py @@ -18,7 +18,7 @@ def __init__(self, table_name: str, partition_statistics_files: List[BaseFile]): super().__init__(table_name) self._partition_statistics_files = partition_statistics_files - def extract_dataframe(self) -> pyspark.sql.DataFrame: + def extract_dataframe(self, include_samples: bool = True) -> pyspark.sql.DataFrame: summaries_df = None samples_df = None for statistics_file in self._partition_statistics_files: @@ -28,14 +28,15 @@ def extract_dataframe(self) -> pyspark.sql.DataFrame: columns = source_df.columns summary_df = self._summarize_partition_statistics(source_df, columns).withColumn("file_path", F.lit(statistics_file.file_path)) - sample_df = self._sample_partitions(source_df, columns).withColumn("file_path", F.lit(statistics_file.file_path)) summaries_df = summary_df if summaries_df is None else summaries_df.unionByName(summary_df, allowMissingColumns=True) - samples_df = sample_df if samples_df is None else samples_df.unionByName(sample_df, allowMissingColumns=True) + if include_samples: + sample_df = self._sample_partitions(source_df, columns).withColumn("file_path", F.lit(statistics_file.file_path)) + samples_df = sample_df if samples_df is None else samples_df.unionByName(sample_df, allowMissingColumns=True) if summaries_df is None: return self._spark.createDataFrame([], StructType([])) - return summaries_df.join(samples_df, on="file_path", how="left") + return summaries_df.join(samples_df, on="file_path", how="left") if include_samples else summaries_df def _read_partition_statistics_file(self, file_path: str) -> pyspark.sql.DataFrame: file_format = PurePosixPath(file_path).suffix.lstrip(".").lower() diff --git a/backend/table_inventory/table_inventory.py b/backend/table_inventory/table_inventory.py index d5057f23..adb17256 100644 --- a/backend/table_inventory/table_inventory.py +++ b/backend/table_inventory/table_inventory.py @@ -304,7 +304,9 @@ def _set_current_table_specs(self): try: current_main_metadata_file = next(metadata_file for metadata_file in self._metadata_files if metadata_file.type == FileType.MAIN_METADATA) - self._current_table_specs = TableMetadataCollector(self._table_name).collect(current_main_metadata_file.file_path) + self._current_table_specs = TableMetadataCollector(self._table_name).collect( + current_main_metadata_file.file_path, self._partition_statistics_files + ) except Exception as e: logger.error( diff --git a/frontend/src/features/docs/content/graph-view.mdx b/frontend/src/features/docs/content/graph-view.mdx index 7ace5ccc..c2bae096 100644 --- a/frontend/src/features/docs/content/graph-view.mdx +++ b/frontend/src/features/docs/content/graph-view.mdx @@ -12,7 +12,7 @@ The Graph view shows all Iceberg metadata objects in your selected range as a di - **Table statistics** - Puffin file with column statistics, such as distinct value counts, for one snapshot. Drawn in amber left of the metadata files and linked from the metadata file that added it. IceGraph shows what the metadata file says about it, without opening it -- **Partition statistics** - Parquet file with per-partition counts, such as records, files, and deletes, for one snapshot. Drawn in silver left of the metadata files and linked from the metadata file that added it. IceGraph opens it and shows how the data is spread across partitions, plus a sample of the most recently updated partitions. The details show `partitions_with_deletes` as a separate count, treating missing or null delete-file counts as zero. Samples use a shared set of fields across the collected statistics files +- **Partition statistics** - Parquet file with per-partition counts, such as records, files, and deletes, for one snapshot. Drawn in silver left of the metadata files and linked from the metadata file that added it. IceGraph opens it and shows how the data is spread across partitions, plus a sample of the most recently updated partitions. The distribution caption **Per partition across all partitions** means it covers every partition in the statistics file. The details show `partitions_with_deletes` as a separate count, treating missing or null delete-file counts as zero. Samples use a shared set of fields across the collected statistics files - **Catalog** - not a file. It is the entry point to the table: the catalog is what points to the current main metadata file, and every read of the table starts there. It is drawn larger and highlighted, labelled with the table name, and sits at the top of the metadata column, with a downward arrow to the main metadata file. Its details show the table properties that do not change between updates: name, UUID, location and format version diff --git a/frontend/src/features/docs/content/metadata-view.mdx b/frontend/src/features/docs/content/metadata-view.mdx index 59b423d4..a0dbb70e 100644 --- a/frontend/src/features/docs/content/metadata-view.mdx +++ b/frontend/src/features/docs/content/metadata-view.mdx @@ -12,6 +12,12 @@ Below the table name, **Table description** shows the string stored in the stand - **Total file size (data + deletes)** includes both file types and is shown in GiB. Hover over the size for exact bytes, or expand **Exact bytes** to view and copy the original byte count. - Missing statistics are shown as **Unavailable**. The page explains when no current snapshot exists or the metadata file does not record it. +### Partition statistics + +When the loaded metadata file references partition statistics for its current snapshot, **Partition statistics** appears directly below **At a glance** in both the Metadata view and the latest metadata page. The section is hidden entirely when no matching statistics exist. + +It shows the partition count, the number of partitions with deletes, and the same minimum, average, and maximum distribution shown for a partition statistics file in the Graph view. **Per partition across all partitions** means the distribution covers every partition in the statistics file. The distribution covers records, data file counts, total data file size, and average data file size per partition. Metadata requests collect only these summaries, without sampling partition rows. If the statistics file cannot be read, its error is shown and missing counts appear as **Unavailable**. + ### Schema and data organization The schema panel scrolls independently. Use **Expand schema** to display its full height. Nested structs, list elements, and map keys and values retain their field IDs and show whether they are required, optional, or unknown. Field comments and identifier fields appear when recorded. Identifier fields do not imply enforced uniqueness or automatic upsert behavior. @@ -26,4 +32,4 @@ Help on Iceberg-specific terms opens on hover, focus, or click. Press Escape to Expand the technical sections for storage paths and identifiers, branches and tags, properties, and delete statistics. JSON-valued properties are formatted and long values can be expanded. -**Advanced: metadata JSON** shows the returned metadata in a scrollable, highlighted block. **Copy metadata JSON** copies exactly the displayed JSON. The backend omits `metadata-log`, `snapshot-log`, `snapshots`, and `statistics` due to size; this export is not the full original metadata file. +**Advanced: metadata JSON** shows the returned metadata in a scrollable, highlighted block, including its `statistics` entries when present. **Copy metadata JSON** copies exactly the displayed JSON. The backend omits `metadata-log`, `snapshot-log`, and `snapshots` due to size; this export is not the full original metadata file. diff --git a/frontend/src/features/metadata/components/MetadataDetails.tsx b/frontend/src/features/metadata/components/MetadataDetails.tsx index de9bbbce..0d68bca1 100644 --- a/frontend/src/features/metadata/components/MetadataDetails.tsx +++ b/frontend/src/features/metadata/components/MetadataDetails.tsx @@ -70,9 +70,8 @@ const MetadataDetails = ({ metadata, onOpenSpec }: MetadataDetailsProps) => {
Reduced metadata: the backend omits metadata-log,{" "}
- snapshot-log, snapshots, and{" "}
- statistics due to size. This is the returned metadata,
- not the complete file.
+ snapshot-log, and snapshots due to size.
+ This is the returned metadata, not the complete file.