When you're building systems that rely on external data, the most insidious failures aren't the ones that crash your service immediately; they're the silent ones. A third-party API subtly changes a field name, adds a new nullable attribute, or worse, sends malformed data that your pipeline happily ingests, only for downstream analytics to break weeks later. If you've ever spent a weekend debugging why a critical report is showing zeroes, you understand the pain of brittle ingestion. This post is for data engineers and architects who need to build Python-based ingestion layers that don't just consume data, but *understand* it, adapt to its evolution, and gracefully handle its imperfections. I'll walk you through architecting a resilient, schema-aware event ingestion pipeline using advanced Pydantic patterns and `asyncio`, ensuring your data quality from the very first byte.
Key Takeaways
- **Schema Versioning with Pydantic**: Leverage multiple Pydantic models to define explicit data contracts and handle schema evolution gracefully, attempting validation against newer versions first.
- **Dynamic Schema Selection**: Implement a robust parsing function that tries different schema versions in a defined order, preventing pipeline failures when data structures change.
- **Dead Letter Queue (DLQ) Pattern**: Isolate and log malformed or unparseable events to a DLQ, ensuring the main pipeline continues processing valid data while providing a clear path for error investigation and remediation.
- **Concurrent Ingestion with `asyncio`**: Utilize `asyncio` and `httpx` to efficiently fetch and process multiple data streams or batches of events concurrently, significantly improving throughput for live feeds.
- **Observable Ingestion Layer**: Build an orchestration layer with logging and hooks for metrics, allowing you to monitor ingestion health, identify data quality issues, and react proactively to upstream changes.
The Problem
The challenge isn't just fetching data; it's fetching data that *constantly shifts*. External RSS feeds, like the GitHub Engineering blog, are a perfect example. While generally stable, they might introduce new fields for authors, change how summaries are structured, or even occasionally publish an entry with missing critical metadata. A naive ingestion script that assumes a fixed schema will either crash or, worse, silently drop valuable information. My goal was to create an ingestion layer that could not only consume these feeds reliably but also anticipate and adapt to these changes without requiring an immediate code deployment for every minor schema tweak.
Data and Sources
For this post, we'll be ingesting the RSS feed from the GitHub Engineering blog. This provides a real-world, semi-structured data source that exhibits the kind of variability we need to address.
* **GitHub Engineering RSS Feed**:
https://github.blog/engineering/feed/
* **`feedparser` library**:
Official PyPI page
* **`Pydantic` library**:
Official documentation
* **`httpx` library**:
Official documentation
Data accessed on 2024-07-29.
Step 1 — The Evolving Data Contract: Defining Baseline and Versioned Schemas
The first sub-problem is how to represent our semi-structured RSS feed data as Python objects while anticipating future schema changes. If we hardcode a single Pydantic model, any new field from the source will either be ignored or cause a validation error if we require it. The solution lies in defining schema versions.
I start with a baseline `FeedEntryV1` model for the most common, expected fields. Then, I introduce `FeedEntryV2` which extends `V1` with new fields, like a more structured `author_detail`. Crucially, both models use `model_extra = 'allow'` (or `extra = 'allow'` in Pydantic v1) to capture any unknown fields into an `extra_data` dictionary. This ensures forward compatibility; if a new field appears that we haven't explicitly modeled yet, it won't break validation and will be preserved for later inspection or processing.
from datetime import datetime
from typing import Any, Dict, List, Optional
from pydantic import BaseModel, Field, ValidationError
# Pydantic v2 models
class AuthorDetail(BaseModel):
name: str
email: Optional[str] = None
href: Optional[str] = None
model_config = {'extra': 'allow'} # Capture any extra author fields
class FeedEntryV1(BaseModel):
title: str
link: str
summary: str
published: datetime = Field(alias='published_parsed') # Renaming for consistency
model_config = {'extra': 'allow'} # Capture any extra top-level fields
class FeedEntryV2(FeedEntryV1):
# V2 introduces a more structured author field
authors: Optional[List[AuthorDetail]] = None
# 'published' is inherited from V1, but we can override or add if needed.
model_config = {'extra': 'allow'} # Ensure V2 also captures extra fields
Here, `FeedEntryV1` captures the essentials. Notice `published: datetime = Field(alias='published_parsed')`. `feedparser` provides `published_parsed` as a `time.struct_time` object, which Pydantic can convert to `datetime` if the field is aliased. `FeedEntryV2` then extends `V1`, adding a `List[AuthorDetail]` to handle multiple authors with structured data. `model_config = {'extra': 'allow'}` is critical; it tells Pydantic to collect any fields not explicitly defined into the `model_extra` attribute (which can be accessed via `entry.model_extra`), preventing validation failures for unexpected but potentially useful data.
Step 2 — Resilient Ingestion: Parsing and Dynamic Schema Selection
With our schema versions defined, the next challenge is reliably parsing raw feed data and selecting the appropriate Pydantic model for validation. A rigid approach would pick one model and fail if the data doesn't conform. My strategy is to try the newest, most comprehensive schema first, then fall back to older versions if validation fails. This ensures we always try to capture the richest data possible.
The `parse_feed_entry` function takes a raw dictionary representation of an entry and attempts to validate it. It prioritizes `FeedEntryV2`, and if that fails with a `ValidationError`, it attempts `FeedEntryV1`. If both fail, it signifies a truly malformed entry that needs special handling. This pattern allows for graceful degradation and ensures that even if a new field is introduced that breaks `V2` for some entries, `V1` can still process the core data.
def parse_feed_entry(entry_data: Dict[str, Any], schema_version: int = 2) -> Optional[BaseModel]:
"""
Attempts to parse a raw feed entry dict against available schema versions,
prioritizing newer versions.
"""
# Convert time.struct_time to datetime for Pydantic
if 'published_parsed' in entry_data and isinstance(entry_data['published_parsed