Skip to content
Merged
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
2 changes: 2 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand Down
39 changes: 24 additions & 15 deletions backend/collectors/collect_partition_statistics.py
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand All @@ -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,
)
Expand Down
55 changes: 52 additions & 3 deletions backend/collectors/collect_table_metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand All @@ -20,24 +22,71 @@ 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 = (
metadata_df.withColumn("current-snapshot", self._current_snapshot_column(metadata_df))
.drop("metadata-log")
.drop("snapshot-log")
.drop("snapshots")
.drop("statistics")
.first()
)

metadata = row.asDict(recursive=True)
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:
Expand Down
9 changes: 5 additions & 4 deletions backend/extractors/partition_statistics_extractor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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()
Expand Down
4 changes: 3 additions & 1 deletion backend/table_inventory/table_inventory.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
2 changes: 1 addition & 1 deletion frontend/src/features/docs/content/graph-view.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
8 changes: 7 additions & 1 deletion frontend/src/features/docs/content/metadata-view.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -70,9 +70,8 @@ const MetadataDetails = ({ metadata, onOpenSpec }: MetadataDetailsProps) => {
<div className={`${UI_EXPANDABLE_BODY_CLASS} space-y-3 px-5 py-4`}>
<p className="text-xs leading-relaxed text-slate-400">
Reduced metadata: the backend omits <code>metadata-log</code>,{" "}
<code>snapshot-log</code>, <code>snapshots</code>, and{" "}
<code>statistics</code> due to size. This is the returned metadata,
not the complete file.
<code>snapshot-log</code>, and <code>snapshots</code> due to size.
This is the returned metadata, not the complete file.
</p>
<MetadataJson
text={JSON.stringify(metadata, null, 2)}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import type { TableMetadata } from "../../table/api/metadataSchemas";
import HelpTerm from "../../../shared/components/HelpTerm";
import MetadataVersion from "./MetadataVersion";
import MetadataSummary from "./MetadataSummary";
import MetadataPartitionStatistics from "./MetadataPartitionStatistics";

interface MetadataOverviewProps {
metadata: TableMetadata;
Expand Down Expand Up @@ -52,6 +53,7 @@ const MetadataOverview = ({
updatedAt={metadata["last-updated-ms"]}
/>
<MetadataSummary metadata={metadata} />
<MetadataPartitionStatistics metadata={metadata} />
</>
);
};
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
import type { TableMetadata } from "../../table/api/metadataSchemas";
import PartitionDistributionTable from "../../table/components/PartitionDistributionTable";
import PanelIssueNotice from "../../../components/PanelIssueNotice";
import { formatCount, integerText } from "../metadataPresentation";

interface MetadataPartitionStatisticsProps {
metadata: TableMetadata;
}

const MetadataPartitionStatistics = ({
metadata,
}: MetadataPartitionStatisticsProps) => {
const statistics = metadata["current-partition-statistics"];
if (!statistics) return null;

return (
<section
aria-labelledby="metadata-partition-statistics-title"
className="space-y-3"
>
<h2
id="metadata-partition-statistics-title"
className="text-base font-semibold text-ink"
>
Partition statistics
</h2>
<PanelIssueNotice type="error">{statistics.errors}</PanelIssueNotice>
<PanelIssueNotice type="warning">{statistics.warnings}</PanelIssueNotice>
<dl className="grid grid-cols-2 divide-x divide-edge rounded-xl border border-edge bg-surface">
<div className="p-5">
<dt className="text-xs text-slate-400">Partitions</dt>
<dd className="mt-2 text-2xl font-semibold text-ink">
{formatCount(integerText(statistics.partitions_count))}
</dd>
</div>
<div className="p-5">
<dt className="text-xs text-slate-400">Partitions with deletes</dt>
<dd className="mt-2 text-2xl font-semibold text-ink">
{formatCount(integerText(statistics.partitions_with_deletes))}
</dd>
</div>
</dl>
<PartitionDistributionTable
partitionDistribution={statistics.partition_distribution}
/>
</section>
);
};

export default MetadataPartitionStatistics;
2 changes: 1 addition & 1 deletion frontend/src/features/table/api/graphCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import {
import { env } from "../../../shared/lib/env";
import { graphDataSchema, type GraphData } from "./graphSchemas";

const GRAPH_CACHE_SCHEMA_VERSION = 3;
const GRAPH_CACHE_SCHEMA_VERSION = 4;
const GRAPH_CACHE_PREFIX = `graph:v${String(GRAPH_CACHE_SCHEMA_VERSION)}:${env.appVersion}:`;
const GRAPH_CACHE_INDEX_PREFIX = "graph-cache-index:";
const GRAPH_CACHE_CLEANUP_KEY = "graph-cache:last-cleanup";
Expand Down
20 changes: 5 additions & 15 deletions frontend/src/features/table/api/graphSchemas.ts
Original file line number Diff line number Diff line change
@@ -1,22 +1,12 @@
import { z } from "zod";
import { icebergIntegerSchema, tableMetadataSchema } from "./metadataSchemas";
import {
icebergIntegerSchema,
partitionDistributionSchema,
tableMetadataSchema,
} from "./metadataSchemas";

export const partitionStatisticsRowSchema = z.record(z.string(), z.unknown());

const statisticValueSchema = z.union([z.number(), z.string()]).nullable();
const minAvgMaxSchema = z.object({
min: statisticValueSchema,
avg: statisticValueSchema,
max: statisticValueSchema,
});

export const partitionDistributionSchema = z.object({
data_record_count: minAvgMaxSchema.optional(),
total_data_file_size_in_bytes: minAvgMaxSchema.optional(),
data_file_count: minAvgMaxSchema.optional(),
average_data_file_size_in_bytes: minAvgMaxSchema.optional(),
});

export const graphDataSchema = z.object({
nodes: z.array(
z.looseObject({
Expand Down
Loading
Loading