Communitygithub.com

SkillMedev/skills

Writes and tunes PySpark jobs - join strategy and broadcast size limits, shuffle-partition sizing, skew diagnosis and salting, UDF avoidance, caching, and output file layout - with concrete size and skew thresholds. Use when someone asks "why is my Spark job slow", "should I broadcast this join", "one task takes forever while the rest finish", "my job OOMs during a join", or is writing a new PySpark ETL job. Do NOT use for Kafka topic, consumer-group, or streaming-pipeline design - use kafka-pipelines instead; do NOT use for single-machine dataframe work that fits in memory - use pandas-expert instead.

skills 是什么?

skills is a Claude Code agent skill that writes and tunes PySpark jobs - join strategy and broadcast size limits, shuffle-partition sizing, skew diagnosis and salting, UDF avoidance, caching, and output file layout - with concrete size and skew thresholds. Use when someone asks "why is my Spark job slow", "should I broadcast this join", "one task takes forever while the rest finish", "my job OOMs during a join", or is writing a new PySpark ETL job. Do NOT use for Kafka topic, consumer-group, or streaming-pipeline design - use kafka-pipelines instead; do NOT use for single-machine dataframe work that fits in memory - use pandas-expert instead.

兼容平台~Claude Code~Codex CLI~Cursor
npx skills add https://github.com/SkillMedev/skills/tree/HEAD/skills/spark-jobs

在你喜欢的 AI 中提问

打开一个已预加载此 Agent Skill 的新对话。

文档

Spark & PySpark

Almost every slow Spark job is slow for one of three reasons: an avoidable shuffle, a skewed key, or rows leaking through Python one at a time. The costly mistake is tuning cluster size before reading the plan - paying for more executors to run the same wasteful DAG faster.

Operating procedure

Step 1: Gather inputs

  1. The job code and, for a slow job, the Spark UI view of the slow stage (task duration distribution, shuffle read/write sizes) and df.explain(True) output.
  2. Input data size and format, and the sizes of both sides of every join.
  3. Cluster shape: executor count, cores, and memory.
  4. Whether the job is new authoring or a tuning pass - for tuning, diagnose from the plan before touching code.

Step 2: Start from a sane session

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = (
    SparkSession.builder
    .appName("etl")
    .config("spark.sql.shuffle.partitions", "200")
    .config("spark.sql.adaptive.enabled", "true")
    .getOrCreate()
)

Enable Adaptive Query Execution (AQE); it coalesces shuffle partitions and switches join strategies at runtime. Then size spark.sql.shuffle.partitions deliberately: target 100-200 MB per shuffle partition, so partitions ≈ shuffle data size ÷ 128 MB, and at least 2-3x the total executor cores so no core idles.

Step 3: Stay in the DataFrame API and out of Python

DataFrame operations go through Catalyst and Tungsten, giving query optimization and off-heap memory. Drop to RDDs only when no DataFrame primitive exists. Python UDFs serialize each row to the Python worker and back, breaking codegen. Order of preference:

  1. Built-in pyspark.sql.functions (e.g. F.regexp_extract, F.when).
  2. Pandas (vectorized) UDFs when Python is unavoidable.
  3. Scalar Python UDFs only as a last resort.
from pyspark.sql.functions import pandas_udf

@pandas_udf("double")
def normalize(s):
    return (s - s.mean()) / s.std()

Step 4: Pick the join strategy by size

  • Broadcast the smaller side when it fits in memory. Spark auto-broadcasts below spark.sql.autoBroadcastJoinThreshold (default 10 MB); explicitly broadcast tables up to a few hundred MB when executor memory allows, and treat ~1 GB as the practical ceiling - broadcast copies the table to every executor and OOMs the job when misjudged.
from pyspark.sql.functions import broadcast
result = large.join(broadcast(small), "key")
  • Filter and project columns before joining to shrink the shuffle; every dropped column and row is shuffle bytes saved.
  • Two large tables → accept the sort-merge join, but check for skew first (Step 5).

Step 5: Diagnose and fix skew

Skew symptoms: the stage is done except for one or two tasks; max task duration exceeds ~3-5x the median; one shuffle-read partition is far larger than the rest in the Spark UI. AQE's skew-join handling splits a partition when it is both over 5x the median partition size and over 256 MB (the defaults) - verify it fired in the plan. When AQE cannot help (e.g. skewed aggregation), salt the hot keys: append a random suffix 0..N to the key on the big side, explode the small side across all N suffixes, join, then strip the salt. Choose N ≈ the skew factor (a key 20x the median gets N=20).

Step 6: Control partitioning and output layout

  • repartition(n, col) does a full shuffle to balance data; use before wide writes.
  • coalesce(n) reduces partitions without a shuffle; use to avoid tiny output files.
  • Target output files of 128 MB-1 GB; thousands of small files punish every downstream reader.
  • Partition output only by low-cardinality columns (date, region - dozens to hundreds of values, never IDs):
df.write.partitionBy("dt").mode("overwrite").parquet("/data/out")

Step 7: Cache only what is reused

Cache only when a DataFrame is used multiple times in the DAG, materialize it, and unpersist when done:

df.cache()
df.count()  # materialize
# ... reuse df ...
df.unpersist()

Caching a once-used DataFrame wastes memory and can evict data that mattered.

Step 8: Verify with the plan

Re-run df.explain(True) and confirm the intended strategy appears (BroadcastHashJoin, pushed filters, no Exchange where one was eliminated). Read columnar formats (Parquet, ORC) so predicate pushdown works, never collect() large data to the driver, and prefer F.col references over string columns when chaining.

Worked example: bad vs good join

Bad - a 2 TB events table joined to a 200 MB dimension, full width, Python UDF for parsing:

result = (events.join(dims, "store_id")          # sort-merge: shuffles all 2 TB
    .withColumn("region", parse_region_udf("meta"))  # row-at-a-time Python
    .filter(F.col("dt") == "2024-06-01"))            # filter AFTER the join

The plan shows an Exchange on both sides (2 TB shuffled), and the UDF blocks codegen.

Good - filter and project first, broadcast the dimension, use a built-in:

result = (events
    .filter(F.col("dt") == "2024-06-01")              # prune to ~30 GB first
    .select("store_id", "amount", "meta")
    .join(broadcast(dims.select("store_id", "region_code")), "store_id")
    .withColumn("region", F.regexp_extract("meta", r"region=(\w+)", 1)))

The plan now shows BroadcastHashJoin with no large-side Exchange: ~30 GB scanned, zero shuffle of the big table, no Python boundary.

Deliverable

Produce the revised job code plus a tuning note stating: the join strategy per join with the size evidence, the shuffle-partition setting with its arithmetic, any skew found (max-vs-median task ratio) and the fix applied, caching decisions, output file-size expectations, and the before/after explain deltas for a tuning pass.

Do NOT

  • Do NOT broadcast a table you have not sized - a misjudged broadcast OOMs every executor at once.
  • Do NOT write a Python UDF before checking pyspark.sql.functions; nearly all string, date, and conditional logic has a built-in.
  • Do NOT collect() or toPandas() a large DataFrame; it funnels the cluster's data through one driver.
  • Do NOT partition output by a high-cardinality column - millions of directories, one row each.
  • Do NOT throw executors at a job whose slow stage is one skewed task; more cores cannot split one partition.
  • Do NOT cache by reflex; cache only DataFrames reused downstream, and unpersist them.

Quality bar

A finished job passes only when: every join names its strategy and the size evidence for it; spark.sql.shuffle.partitions is derived from data size, not left at a habit value; the Spark UI shows no task above ~3x the median duration in shuffle stages (or the skew is explained and fixed); no scalar Python UDF survives where a built-in exists; output files land in the 128 MB-1 GB band; and the final explain confirms the intended plan.

Individual skills in this repo

This repo contains 11 individual skills — each has its own dedicated page.

SkillMedev/skills

Designs REST API surfaces - resource naming, HTTP method and status-code semantics, error shapes, pagination, and filtering - and delivers an endpoint spec a consumer can build against without asking questions. Use when someone asks "how should I name this endpoint", "what status code should this return", "should this be PUT or PATCH", "how do I paginate this list", or is reviewing an API before it ships to external consumers. Do NOT use for planning breaking-change rollouts and deprecation windows - use api-versioning-strategist instead; for GraphQL type and resolver design - use graphql-schema instead; for generating client SDKs from an existing spec - use api-client-generator instead; for designing inbound webhook endpoints - use webhook-receiver-hardener instead.

SkillMedev/skills

Turns data and charts into a decision-driving narrative structured as headline finding, trend, implication, and recommended action - with finding-led chart titles, context for every number, annotation guidance, and honest flags on any conclusion the data cannot support. Use when someone says "turn these numbers into a story", "what's the takeaway from this data", "help me present these results to leadership", or has charts but no narrative. Do NOT use for compressing a long document into a one-pager - use executive-summary instead - or for running the analysis that produces the findings - use eda-playbook instead.

SkillMedev/skills

Use when a task needs live or historical money data - "convert USD to EUR", "current/past exchange rate", "FX rate on this date / over this range", or "current price of Bitcoin/Ethereum, market cap, 24h change". Frankfurter (ECB reference rates, no key) is the FX default; CoinGecko's free keyless tier covers crypto. Do NOT use for stock quotes or equities - no keyless stock API survives verification, say so instead of guessing; do NOT use for country economic indicators like GDP or inflation series - use government-open-data instead; if the request is a vague "I need live data", route through public-data-api-picker.

SkillMedev/skills

Builds a driver-based FP&A operating model linking business inputs to P&L, balance sheet, and cash flow outputs. Use when building an annual plan, preparing investor materials, running scenario analysis, or stress-testing the business.

SkillMedev/skills

Use when a task needs live geographic lookups - "geocode this address", "what's at these coordinates" (reverse geocoding), "lat/lon for this city", "which country/state is this ZIP or postal code in", or "country facts: capital, currency, population, flag". Nominatim (OpenStreetMap) is the geocoding default; Zippopotam for postal codes; APICountries for country facts. All keyless. Do NOT use for weather at a location - use weather-climate instead; do NOT use for country-level statistics over time (GDP, population trends) - use government-open-data instead; if the request is a vague "I need live data", route through public-data-api-picker.

SkillMedev/skills

Runs the full Getting Things Done loop - capture, clarify, organize, reflect, engage - building a trusted system of context lists, a projects list with defined next actions, and a weekly review habit. Use when someone says "I'm overwhelmed and things are slipping through the cracks", "set up GTD for me", "help me do a brain dump and organize it", or "my to-do list is a mess". Do NOT use for just running the weekly review ritual itself - use weekly-review instead - or for clearing an email backlog - use inbox-zero.

SkillMedev/skills

Processes any email backlog to zero using the 4Ds - Delete, Delegate, Defer, Do - with a mass-archive strategy for the obvious, a touch-each-email-once discipline, and a keep-it-clear system of batched processing windows, ruthless unsubscribing, filters, and a minimal folder setup. Use when someone says "I have 5,000 unread emails", "help me get to inbox zero", "email is eating my whole day", or treats their inbox as a to-do list. Do NOT use for drafting the reply emails themselves or prioritization rules for an ongoing support queue - use email-triage instead - or for protecting focus time around the email windows - use deep-work-planner instead.

SkillMedev/skills

Runs structured coaching sessions using values clarification and the GROW model, ending every session with one committed action, a deadline, and an if-then plan for the likely obstacle. Use when someone says "I feel stuck in my life", "help me figure out what I want", "hold me accountable to my goals", or "coach me through this decision". Do NOT use for building a stress toolkit - use stress-management instead - or a journaling practice - use journal-framework; for a standing goal-tracking system, use goals-accountability. Coaching, not therapy: signs of clinical distress route to a licensed professional.

SkillMedev/skills

Classifies incident severity (SEV1-4) using impact, scope, and urgency signals and decides who to page. Use when an alert fires or a report comes in and a severity call must be made quickly.

SkillMedev/skills

Use the Skill Me catalog from inside any conversation - discover, install, and manage Claude skills through the Skill Me MCP, and load installed skills automatically each session.

SkillMedev/skills

Builds clean, performant, accessible SwiftUI views with correct state ownership, scoped invalidation, and smooth list scrolling, and reviews existing SwiftUI code against a concrete frame-time and re-render budget. Use when someone asks "why does my SwiftUI list stutter", "should this be @State or @Observable", "my whole screen re-renders when one row changes", "how do I animate this transition", or wants a SwiftUI view built or refactored. Do NOT use for cross-platform React Native apps - use react-native-pro instead; do NOT use for Flutter widget trees - use flutter-widget-architect instead; do NOT use for Android Compose UIs - use jetpack-compose-builder instead.

相关技能