Loading skill
Install any skill in seconds. Free to start, no credit card required.
Get Started Free →Process large-scale data with Apache Spark. Use when a user asks to process big data, run distributed computations, build ETL pipelines, perform data analysis at scale, or use PySpark for data engineering.
| Test case | Without → With | Effect | Δ tokens | Δ turns |
|---|---|---|---|---|
| case-01 | ✗→✓ | ▲ Improved | 0% | 0% |
| case-02 | ✗→✓ | ▲ Improved | 33% | 0% |
| case-21 | ✓→✓ | = Same ✓ | 13% | 0% |
| case-12 | ✓→✓ | = Same ✓ | 112% | 0% |
| case-04 | ✓→✓ | = Same ✓ | 8% | 0% |
Apache Spark is the standard for distributed data processing. It handles batch processing, streaming, SQL, machine learning, and graph processing. PySpark provides a Python API. Runs on standalone clusters, YARN, Kubernetes, or managed services (Databricks, EMR, Dataproc).
bashpip install pyspark
python# etl/process.py — PySpark data processing from pyspark.sql import SparkSession from pyspark.sql import functions as F spark = SparkSession.builder \ .appName("DataPipeline") \ .config("spark.sql.adaptive.enabled", "true") \ .getOrCreate() # Read data df = spark.read.parquet("s3://bucket/raw/events/") # Transform processed = (df .filter(F.col("event_type").isin(["purchase", "signup"])) .withColumn("date", F.to_date("timestamp")) .withColumn("revenue", F.col("amount") * F.col("quantity")) .groupBy("date", "event_type") .agg( F.count("*").alias("event_count"), F.sum("revenue").alias("total_revenue"), F.countDistinct("user_id").alias("unique_users"), ) .orderBy("date") ) # Write results processed.write \ .mode("overwrite") \ .partitionBy("date") \ .parquet("s3://bucket/processed/daily_metrics/")
python# Register as SQL table df.createOrReplaceTempView("events") result = spark.sql(""" SELECT date_trunc('month', timestamp) as month, COUNT(DISTINCT user_id) as monthly_active_users, SUM(CASE WHEN event_type = 'purchase' THEN amount ELSE 0 END) as revenue FROM events WHERE timestamp >= '2025-01-01' GROUP BY 1 ORDER BY 1 """) result.show()
python# Real-time processing from Kafka stream = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "kafka:9092") \ .option("subscribe", "events") \ .load() parsed = stream.select( F.from_json(F.col("value").cast("string"), schema).alias("data") ).select("data.*") query = parsed \ .groupBy(F.window("timestamp", "5 minutes"), "event_type") \ .count() \ .writeStream \ .outputMode("update") \ .format("console") \ .start()
Other measured skills in the registry, with their headline benchmark lift.