You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

优化Databricks上PySpark的HTTP请求性能(Requests模块)

Optimizing HTTP Status Code Fetching in PySpark on Databricks

Hey Ashley, 14 hours for 25k HTTP requests definitely sounds way too slow—let’s figure out what’s going wrong and how to speed this up. Most likely, the issue comes from how you’re handling the requests combined with suboptimal Spark configuration. Here’s a breakdown of common pitfalls and fixes:

Common Reasons for Slow Performance

  • Synchronous HTTP calls in Python UDFs: If you’re using a library like requests in a standard PySpark UDF, each request blocks the worker thread until it completes. Plus, Python UDFs have overhead from serializing data between JVM and Python processes, which adds up for 25k individual calls.
  • Poor parallelism configuration: If your DataFrame has too few partitions, each worker is stuck handling hundreds/thousands of requests sequentially. Too many partitions, and you waste resources on task startup overhead.
  • No connection pooling or timeouts: Without reusing TCP connections (via a connection pool), each request has to go through the full TCP handshake. Missing timeouts can leave requests hanging indefinitely if a server doesn’t respond.

Performance Optimization Strategies

1. Use Asynchronous HTTP Requests per Partition

Instead of handling one request at a time, process batches of URLs in each Spark partition using an async HTTP client like aiohttp. This lets each worker handle multiple requests concurrently, drastically reducing total time.

Here’s a sample implementation:

import aiohttp
import asyncio
from pyspark.sql import SparkSession

# Async function to fetch a single URL's status
async def fetch_status(session, url):
    try:
        # Set reasonable timeouts to avoid hanging requests
        async with session.get(
            url,
            timeout=aiohttp.ClientTimeout(total=10),
            headers={"User-Agent": "Mozilla/5.0"}  # Mimic a browser to avoid blocks
        ) as response:
            return (url, response.status)
    except Exception as e:
        # Return error details instead of failing the entire partition
        return (url, f"Error: {str(e)}")

# Async function to process all URLs in a single partition
async def process_partition(urls):
    # Use a TCP connector with connection pooling
    async with aiohttp.ClientSession(
        connector=aiohttp.TCPConnector(limit=20)  # Limit concurrent requests per partition
    ) as session:
        tasks = [fetch_status(session, url) for url in urls]
        return await asyncio.gather(*tasks)

# Wrapper to run async code in Spark's Python worker
def partition_worker(partition):
    # Spark's Python workers use multiprocessing, so create a new event loop per process
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    urls = [row.url for row in partition]
    results = loop.run_until_complete(process_partition(urls))
    loop.close()
    return results

# Example usage with your DataFrame
spark = SparkSession.builder.appName("URLStatusChecker").getOrCreate()
url_df = spark.read.table("your_url_table")  # Replace with your data source

# Repartition to balance work: aim for 50-150 URLs per partition
# For 25k URLs, 200-500 partitions is a good starting point
optimized_df = url_df.repartition(300)

# Process each partition with async requests
status_results = optimized_df.rdd.mapPartitions(partition_worker).toDF(["url", "status"])

# Save results
status_results.write.mode("overwrite").saveAsTable("url_status_results")

2. Tune Spark Parallelism

  • Adjust partition count: As a rule of thumb, set partitions to 2-3 times the number of CPU cores in your cluster. For example, if you have a cluster with 10 workers each with 4 cores, aim for 80-120 partitions. This ensures full utilization of your cluster resources.
  • Avoid skewed partitions: If some URLs take much longer to process (e.g., slow servers), use salting or split skewed partitions to prevent straggler tasks.

3. Use Connection Pooling & Timeouts

  • Connection pooling (like TCPConnector in aiohttp) reuses existing TCP connections, cutting down on handshake overhead.
  • Always set total and connect timeouts to prevent requests from hanging. Start with 5-10 seconds and adjust based on your target servers’ response times.

4. Consider Scala UDFs for Even Better Performance

Python UDFs have inherent overhead from cross-language serialization. If you’re comfortable with Scala, using a JVM-based async HTTP client (like AsyncHttpClient) can give you another performance boost. Scala runs directly in the Spark JVM, eliminating the Python process overhead.

5. Databricks-Specific Tweaks

  • Cluster configuration: Use worker instances with sufficient network bandwidth (e.g., m5.xlarge or higher) to handle concurrent requests. Avoid under-provisioned clusters like the Community Edition’s default small instances.
  • Runtime version: Use a recent Databricks Runtime (11.3 LTS or newer) to leverage the latest Spark optimizations and bug fixes.

Final Notes

  • Be mindful of rate limits: Too many concurrent requests might get your IP blocked by target servers. Adjust the limit in TCPConnector based on the server’s rate policies.
  • Handle retries for transient errors: Add a retry mechanism (e.g., using tenacity with async support) for 5xx errors or timeouts to improve result accuracy without reprocessing all URLs.

With these changes, you should see your runtime drop from hours to minutes for 25k URLs.

内容的提问来源于stack exchange,提问作者Ashley O

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 07:57:56