▸case-01 I have a PySpark ETL script processing 500GB of log data across 100 partitions with severe data skew that runs slowly during shuffle joins. Our cluster runs 20 workers with 8GB memory per executor. Please review the pipeline and return the refactored code, explicit tuning recommendations, projected percentage improvements for runtime, memory/CPU usage, and cost, alongside the required Spark configuration settings. | fail→pass | 23,293 | 10,413 | -55% | 1 | 1 | 0% | 4,299 | 2,534 | -41% | 0 | 0 | — |
▸case-02 Our nightly Spark batch job reading Delta tables is hitting out-of-memory errors on executors during aggregation steps. Here is the transformation logic along with cluster metrics. Could you provide optimized Spark code, an array of actionable recommendations, expected percentage changes for execution time, resource consumption, and cost, as well as a JSON object of updated cluster configurations? | fail→pass | 17,884 | 10,735 | -40% | 1 | 1 | 0% | 3,171 | 2,530 | -20% | 0 | 0 | — |
▸case-03 We are migrating a streaming aggregation job and want to resolve severe shuffle bottlenecks due to unoptimized joins. Attached is the PySpark source along with system resource details. Please evaluate the code and output an updated code snippet, strategic advice list, estimated percentage gains across timing and resource expenses, and the target configuration dictionary. | fail→pass | 15,447 | 15,276 | -1% | 1 | 1 | 0% | 2,837 | 3,583 | +26% | 0 | 0 | — |
▸case-04 We have a PySpark Structured Streaming job consuming from Kafka and performing windowed aggregations. We need an optimization evaluation structured as a JSON payload containing optimized code, tuning recommendations, expected improvement metrics (execution time, resource usage, cost), and configuration updates. | fail→pass | 17,997 | 9,780 | -46% | 1 | 1 | 0% | 3,471 | 2,342 | -33% | 0 | 0 | — |
▸case-05 We are building a feature store generation pipeline in PySpark joining daily clickstream data with user profile tables. We require an optimization analysis formatted strictly as JSON matching schema fields for optimized code, recommendation list, expected percentage gains (executionTime, resourceUsage, cost), and config dictionary. | fail→pass | 12,845 | 11,450 | -11% | 1 | 1 | 0% | 2,704 | 2,644 | -2% | 0 | 0 | — |
▸case-06 We are migrating a legacy SQL query pipeline to PySpark DataFrames. The job processes 2TB of partitioned transactional records. Analyze the job and respond with a JSON object containing optimizedCode, recommendations, expectedImprovement, and configChanges. | pass→pass | 19,273 | 11,276 | -41% | 1 | 1 | 0% | 3,186 | 2,669 | -16% | 0 | 0 | — |
▸case-07 In a PySpark join between a 1TB skew-heavy event log dataset and a 50MB dimension table, executor tasks processing key 'user_null' take 45 minutes while others finish in 10 seconds. Format the response as JSON with optimizedCode, recommendations, expectedImprovement, and configChanges. How should the code handle key skew? | pass→pass | 10,734 | 10,399 | -3% | 1 | 1 | 0% | 2,019 | 2,497 | +24% | 0 | 0 | — |
▸case-08 A PySpark ETL job joins a 150MB dimension DataFrame with a 500GB fact table. Spark is spilling to disk during shuffle hash joins because auto-broadcast is disabled. Respond in JSON schema format (optimizedCode, recommendations, expectedImprovement, configChanges) with specific configuration key-value pairs for autoBroadcastJoinThreshold. | pass→pass | 8,064 | 9,035 | +12% | 1 | 1 | 0% | 1,890 | 2,429 | +29% | 0 | 0 | — |
▸case-09 A PySpark job outputs 50,000 tiny 200KB files after filtering a 100GB dataset down to 1GB. The developer used `repartition(1)` causing a massive single-executor shuffle bottleneck. Provide the optimized JSON response containing optimizedCode, recommendations, expectedImprovement, and configChanges. | pass→pass | 10,634 | 9,095 | -14% | 1 | 1 | 0% | 2,157 | 2,377 | +10% | 0 | 0 | — |
▸case-10 A PySpark Pandas UDF processing large batches fails repeatedly with container memory limit exceeded errors on YARN. Executor memory is set to 16g. Provide a JSON output (optimizedCode, recommendations, expectedImprovement, configChanges) with the required Spark configuration to allocate off-heap memory. | pass→pass | 10,856 | 8,978 | -17% | 1 | 1 | 0% | 2,041 | 2,157 | +6% | 0 | 0 | — |
▸case-11 A complex PySpark DataFrame transformation has nested subqueries that Spark's Catalyst engine evaluates sub-optimally. Provide a JSON response (optimizedCode, recommendations, expectedImprovement, configChanges) demonstrating explicit Catalyst join hints. | pass→pass | 12,907 | 9,533 | -26% | 1 | 1 | 0% | 2,391 | 2,219 | -7% | 0 | 0 | — |
▸case-12 A 5GB PySpark batch job runs with default Spark shuffle settings (200 partitions), resulting in 200 tiny tasks processing 25MB each and incurring high scheduling overhead. Return a JSON object with optimizedCode, recommendations, expectedImprovement, and configChanges detailing AQE configuration. | pass→pass | 9,983 | 8,724 | -13% | 1 | 1 | 0% | 2,140 | 2,204 | +3% | 0 | 0 | — |
▸case-13 A PySpark ML feature engineering script evaluates the same filtered DataFrame five times across different aggregations without persisting it, causing repeated re-computations of expensive parses. Provide the JSON response (optimizedCode, recommendations, expectedImprovement, configChanges) correcting this pattern. | pass→pass | 8,899 | 8,465 | -5% | 1 | 1 | 0% | 1,800 | 2,207 | +23% | 0 | 0 | — |
▸case-14 A Spark job on YARN is configured with 1 executor per node having 32 cores and 64GB RAM, suffering from heavy garbage collection pauses and HDFS throughput bottlenecks. Provide a JSON response (optimizedCode, recommendations, expectedImprovement, configChanges) optimizing executor allocation. | pass→pass | 10,359 | 9,080 | -12% | 1 | 1 | 0% | 2,105 | 2,138 | +2% | 0 | 0 | — |
▸case-15 An intermittent Spark batch processing pipeline running on a Kubernetes cluster wastes idle cloud resources when waiting for upstream S3 data partitions. Output a JSON object (optimizedCode, recommendations, expectedImprovement, configChanges) enabling dynamic allocation. | pass→pass | 13,065 | 10,731 | -18% | 1 | 1 | 0% | 2,643 | 2,462 | -7% | 0 | 0 | — |
▸case-16 A PySpark star-schema query joins a multi-terabyte partitioned fact table with a small date dimension table, but Spark scans all partitions of the fact table despite filtering on dimension dates. Provide a JSON response (optimizedCode, recommendations, expectedImprovement, configChanges) enabling dynamic partition pruning. | pass→pass | 10,170 | 8,979 | -12% | 1 | 1 | 0% | 1,949 | 2,036 | +4% | 0 | 0 | — |
▸case-17 Two 500GB PySpark DataFrames are repeatedly joined on `customer_id` in hourly batch jobs, generating massive network shuffles every run. Output a JSON response (optimizedCode, recommendations, expectedImprovement, configChanges) recommending a storage strategy to eliminate join shuffles. | pass→pass | 10,360 | 12,570 | +21% | 1 | 1 | 0% | 1,903 | 2,805 | +47% | 0 | 0 | — |
▸case-18 A PySpark orderBy query on a 200GB dataset causes frequent executor disk spills in Spark UI stage metrics. Provide a JSON response (optimizedCode, recommendations, expectedImprovement, configChanges) addressing executor memory fraction allocation. | pass→pass | 10,787 | 10,395 | -4% | 1 | 1 | 0% | 2,143 | 2,548 | +19% | 0 | 0 | — |
▸case-19 A data engineer wrote a PySpark job using `.toPandas()` on a 50GB DataFrame to compute group aggregations, resulting in driver OutOfMemory exceptions. Return a JSON response (optimizedCode, recommendations, expectedImprovement, configChanges) replacing driver-side collect with distributed Spark operations. | pass→pass | 8,151 | 8,445 | +4% | 1 | 1 | 0% | 1,637 | 2,122 | +30% | 0 | 0 | — |
▸case-20 How do I define an explicit StructType schema with IntegerType and StringType fields for reading a CSV in PySpark? | pass→pass | 5,917 | 4,293 | -27% | 1 | 1 | 0% | 1,267 | 1,163 | -8% | 0 | 0 | — |
▸case-21 How do I compute a dense rank over partitions ordered by timestamp in a PySpark DataFrame using Window functions? | pass→pass | 8,460 | 7,891 | -7% | 1 | 1 | 0% | 1,918 | 1,982 | +3% | 0 | 0 | — |
▸case-22 In PySpark, how do I filter a DataFrame column of array type to keep rows where the array contains the value 'active'? | pass→pass | 6,919 | 5,173 | -25% | 1 | 1 | 0% | 1,331 | 1,367 | +3% | 0 | 0 | — |