In the unpredictable world of production data pipelines, it's not a matter of if you'll encounter data corruption, erroneous updates, or a critical need to inspect data as it was yesterday, but when. Traditional data lakes, often a collection of disparate files, turn these incidents into complex, time-consuming recovery efforts, sometimes compromising data integrity permanently. For experienced data engineers, this post will demystify how Apache Iceberg's core architecture, specifically its snapshot management and time travel features, offers an elegant solution to these challenges, enabling transparent data auditing and swift, confident recovery from even the most severe data incidents.
Key Takeaways
- Establishing a local PyIceberg environment with DuckDB allows for rapid prototyping and testing of Iceberg features against dynamic API data.
- Every data modification (append, overwrite, delete) in an Iceberg table automatically creates a new, immutable snapshot, preserving the table's state at that exact moment.
- Iceberg's time travel enables querying historical data states using either a specific `snapshot_id` or a `timestamp`, providing a powerful tool for auditing and debugging.
- Data rollbacks can revert an Iceberg table to any previous, known-good snapshot, offering a robust mechanism for incident recovery from bad writes or schema changes.
- Understanding snapshot metadata and implementing sensible expiration policies are critical for managing storage costs and maintaining query performance in production Iceberg deployments.
The Challenge of Mutable Data in Production
Imagine a scenario: your critical upstream API feed unexpectedly starts sending malformed data, or a bug in your transformation logic introduces a subtle but pervasive error into your production dataset. Weeks later, a downstream report surfaces the issue. Without a robust mechanism to inspect historical states or revert to a known-good version, you're left with a forensic nightmare, patching data, or even worse, admitting to data loss. This isn't just about backups; it's about having an inherent, immutable ledger of every change, ready for instantaneous query or rollback. This is precisely where Apache Iceberg shines, transforming your data lake into an auditable, recoverable data lakehouse.
Data and Sources
For this exploration, we'll be interacting with the public GitHub API, specifically querying information about the python/cpython repository. This provides a real, dynamic data source that can change over time, allowing us to simulate data evolution. We'll use PyIceberg for Pythonic interaction with Iceberg tables and DuckDB as our local, in-process SQL query engine, providing a lightweight yet powerful setup.
- GitHub API documentation: https://docs.github.com/en/rest/repos/repos
- Apache Iceberg documentation: https://iceberg.apache.org/docs/latest/
- PyIceberg documentation: https://py.iceberg.apache.org/
- DuckDB documentation: https://duckdb.org/docs/
- Source code for the complete script will be available at: https://github.com/mausamadhikari/iceberg-time-travel-demo (Note: This link is illustrative and would point to a real GitHub repo).
Data accessed on 2024-07-25
Step 1 — Initializing the Iceberg Table and First Snapshot
The first step in building an auditable data system is to establish a foundation: creating an Iceberg table and ingesting the initial dataset. This action inherently creates the very first snapshot of our data, marking the beginning of our immutable ledger. This sub-problem is about setting up the local catalog and writing the first version of our external API data.
I'll use a local SqlCatalog backed by DuckDB for simplicity, which stores Iceberg metadata in a local SQLite file. After fetching the initial CPython repo data, I'll convert it into a Pandas DataFrame and then to an Arrow Table, which PyIceberg can write directly. The schema is inferred from the initial data, but defining it explicitly with pyarrow.schema provides more control and robustness.
from pyiceberg.catalog import SqlCatalog
from pyiceberg.schema import Schema
from pyiceberg.types import LongType, StringType, TimestampType
import pandas as pd
import pyarrow as pa
import requests
from datetime import datetime
import time
CATALOG_NAME = "local_catalog"
DB_FILE = "iceberg_catalog.db"
TABLE_NAME = "github_cpython_repo"
REPO_API_URL = "https://api.github.com/repos/python/cpython"
def create_initial_table(catalog):
print("Fetching initial GitHub CPython repo data...")
try:
response = requests.get(REPO_API_URL, timeout=10)
response.raise_for_status() # Raise an HTTPError for bad responses (4xx or 5xx)
data = response.json()
except requests.exceptions.RequestException as e:
print(f"Error fetching data from GitHub API: {e}")
return None
# Select relevant fields and add a timestamp
repo_data = {
"id": data["id"],
"node_id": data["node_id"],
"name": data["name"],
"full_name": data["full_name"],
"description": data["description"],
"stars": data["stargazers_count"],
"forks": data["forks_count"],
"open_issues": data["open_issues_count"],
"fetched_at": datetime.now(),
}
df = pd.DataFrame([repo_data])
# Define the Iceberg schema explicitly
iceberg_schema = Schema(
LongType(1, "id", required=True),
StringType(2, "node_id", required=True),
StringType(3, "name", required=True),
StringType(4, "full_name", required=True),
StringType(5, "description"),
LongType(6, "stars"),
LongType(7, "forks"),
LongType(8, "open_issues"),
TimestampType(9, "fetched_at"),
)
# Convert pandas DataFrame to PyArrow Table
arrow_table = pa.Table.from_pandas(df, schema=iceberg_schema.as_arrow(), preserve_index=False)
print(f"Creating Iceberg table '{TABLE_NAME}' and writing initial data...")
table = catalog.create_table(
identifier=TABLE_NAME,
schema=iceberg_schema,
location=f"./{TABLE_NAME}", # Store table data in a local directory
properties={"format-version": "2"}
)
table.append(arrow_table)
print(f"Table '{TABLE_NAME}' created with initial snapshot {table.current_snapshot.snapshot_id}")
return table
This snippet initializes the catalog, fetches data, structures it, and writes it to a new Iceberg table. The `table.append()` call is crucial here; it not only writes the data but also commits a new snapshot to the table's metadata, effectively capturing this initial state.
Step 2 — Simulating Data Evolution and Capturing New Snapshots
To truly appreciate Iceberg's auditability, we need to see how it handles changes. This step addresses how dynamic changes from an external source are captured as distinct, new snapshots within the Iceberg table, building a historical record without overwriting past data. Since the GitHub API might not change rapidly enough for a short demo, I'll simulate some minor variations in the fetched data to ensure we generate different snapshots.
I'll fetch the data again, introduce a programmatic change to the 'stars' and 'forks' count (e.g., adding a small random number), and then append this new version to the table. Each `append` operation automatically triggers the creation of a new snapshot, linking it to the previous one and preserving the full history.
import random
def append_new_data(table, iteration):
print(f"\nFetching GitHub CPython repo data for iteration {iteration}...")
try:
response = requests.get(REPO_API_URL, timeout=10)