— No reviews yet
0 installs
21 views
0.0% view→install
Install
$ agentstack add skill-chandrudp29-skillhub-data-pipeline ✓ scanned · ✓ verified — works with Claude Code, Cursor, and more.
Security review
✓ PassedNo issues found. Passed automated security review. · v0.1.0 How review works →
- ✓ Prompt-injection patterns
- ✓ Secret / credential exfiltration
- ✓ Dangerous shell & filesystem operations
- ✓ Untrusted network calls
- ✓ Known-malicious package signatures
What it can access
- ✓ Network access No
- ✓ Filesystem access No
- ✓ Shell / process execution No
- ✓ Environment & secrets No
- ✓ Dynamic code execution No
From automated source analysis of v0.1.0. “Used” means the capability is present in the source — more access means more to trust, not that it’s unsafe.
Are you the author of Data Pipeline? Claim this listing to set pricing, connect Stripe payouts, and keep 70% of every sale.
Sign up to claimAbout
When to Use
Apply when building, reviewing, or debugging data pipelines, ETL jobs, or data orchestration workflows.
Core Rules
- All pipelines must be idempotent — running twice gives the same result as running once
- Prefer ELT over ETL — load raw first, transform in the warehouse (cheaper to re-transform than re-extract)
- Always use incremental loads over full loads for tables > 100K rows
- Fail fast and loudly — a silent partial load is worse than a visible failure
- Test transformations with fixed fixtures, not prod data samples
Idempotency Patterns
# ❌ Non-idempotent — double run = double rows
def load_events(conn, events):
for event in events:
conn.execute("INSERT INTO events VALUES (?)", event)
# ✓ Idempotent — delete-then-insert by partition
def load_events(conn, events, date: str):
conn.execute("DELETE FROM events WHERE date = ?", date)
conn.executemany("INSERT INTO events VALUES (?)", events)
# ✓ UPSERT — idempotent by natural key
def upsert_users(conn, users):
conn.executemany("""
INSERT INTO users (id, email, name, updated_at)
VALUES (?, ?, ?, ?)
ON CONFLICT (id) DO UPDATE SET
email = excluded.email,
name = excluded.name,
updated_at = excluded.updated_at
""", users)
Incremental Load Pattern
from datetime import datetime, timedelta
def get_last_watermark(conn, table: str) -> datetime:
row = conn.execute(
"SELECT MAX(watermark) FROM pipeline_state WHERE table_name = ?", table
).fetchone()
return row[0] or datetime(2020, 1, 1)
def update_watermark(conn, table: str, watermark: datetime):
conn.execute("""
INSERT INTO pipeline_state (table_name, watermark)
VALUES (?, ?)
ON CONFLICT (table_name) DO UPDATE SET watermark = excluded.watermark
""", (table, watermark))
def load_incremental(source_conn, dest_conn, table: str):
last = get_last_watermark(dest_conn, table)
new_watermark = datetime.utcnow()
rows = source_conn.execute(
f"SELECT * FROM {table} WHERE updated_at > ?", last
).fetchall()
upsert_rows(dest_conn, table, rows)
update_watermark(dest_conn, table, new_watermark)
Airflow DAG Pattern
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'data-team',
'retries': 3,
'retry_delay': timedelta(minutes=5),
'email_on_failure': True,
}
with DAG(
dag_id='user_events_pipeline',
default_args=default_args,
schedule='0 */6 * * *', # every 6 hours
start_date=datetime(2024, 1, 1),
catchup=False, # don't backfill missed runs on first deploy
tags=['events', 'incremental'],
) as dag:
extract = PythonOperator(
task_id='extract_events',
python_callable=extract_events_from_source,
)
validate = PythonOperator(
task_id='validate_events',
python_callable=run_data_quality_checks,
)
load = PythonOperator(
task_id='load_to_warehouse',
python_callable=load_events_incremental,
)
extract >> validate >> load # explicit dependency chain
Data Quality Checks
# Run these before loading to prod — fail the pipeline, not silently skip
def validate(df):
checks = [
(df['user_id'].notna().all(), "user_id has nulls"),
(df['event_time'].notna().all(), "event_time has nulls"),
(df['event_type'].isin(VALID_TYPES).all(), "unknown event_type values"),
(len(df) > 0, "empty result — no events extracted"),
((df['amount'] >= 0).all(), "negative amounts"),
]
failures = [msg for passed, msg in checks if not passed]
if failures:
raise ValueError(f"Data quality failed: {', '.join(failures)}")
Pipeline Monitoring
- Emit row counts at each stage — alert on drops > 20% vs same window last week
- Track pipeline duration — alert if > 2× historical average
- Dead-letter queue for rows that fail transformation — don't silently drop
- Log: extractcount, transformcount, loadcount, rejectedcount, duration_seconds
- Idempotency test: run pipeline twice on same data, result must be identical
Source & license
This open-source skill is cataloged on AgentStack and links to its original source — we do not rehost the code.
- Author: chandrudp29
- Source: chandrudp29/skillhub
- License: MIT
- Homepage: https://pypi.org/project/skillhub-ai/
Install and usage instructions live in the source repository linked above.
Reviews
No reviews yet — be the first.
Write a review
Versions
- v0.1.0 Imported from the upstream source.