PySpark地理距离计算的向量化实现优化求助
Hey there! I totally get the frustration with slow non-vectorized loops for this kind of geospatial task—let's fix that with vectorized operations using NumPy and Pandas, which will drastically speed things up. Here's a step-by-step breakdown:
Key Idea
Instead of iterating over each row pair (which is glacial in Python), we'll use NumPy broadcasting to compute all pairwise geospatial distances at once, then filter and sum the relevant income values. We'll use the Haversine formula to calculate accurate spherical distances (since latitude/longitude are on a sphere, Euclidean distance won't work here).
Step 1: Vectorized Haversine Function
First, define a vectorized version of the Haversine formula that can handle array inputs (this leverages NumPy's C-level optimizations):
import numpy as np import pandas as pd def haversine(lat1, lon1, lat2, lon2): # Convert degrees to radians (required for trigonometric functions) lat1_rad = np.radians(lat1) lon1_rad = np.radians(lon1) lat2_rad = np.radians(lat2) lon2_rad = np.radians(lon2) # Haversine formula calculations dlat = lat2_rad - lat1_rad dlon = lon2_rad - lon1_rad a = np.sin(dlat/2)**2 + np.cos(lat1_rad) * np.cos(lat2_rad) * np.sin(dlon/2)**2 c = 2 * np.arcsin(np.sqrt(a)) # Earth radius in meters (6,371,000 m) earth_radius = 6371000 return c * earth_radius
Step 2: Prepare Data as NumPy Arrays
Extract the coordinate columns and income from your DataFrames as NumPy arrays—this is where the vectorized magic starts:
# Extract coordinates from df (shape: [number of df rows, 2]) df_coords = df[["latitude", "longitude"]].to_numpy() # Extract coordinates and income from cent (coords shape: [number of cent rows, 2]; income shape: [number of cent rows]) cent_coords = cent[["latitude", "longitude"]].to_numpy() cent_incomes = cent["income"].to_numpy()
Step 3: Compute Pairwise Distances with Broadcasting
Use NumPy broadcasting to create a distance matrix where each element distances[i,j] is the distance between df.iloc[i] and cent.iloc[j]:
# Reshape df coordinates to enable broadcasting with cent coordinates lat_df, lon_df = df_coords[:, 0, np.newaxis], df_coords[:, 1, np.newaxis] lat_cent, lon_cent = cent_coords[:, 0], cent_coords[:, 1] # Calculate all pairwise distances (result is a 2D array: [df rows x cent rows]) distances = haversine(lat_df, lon_df, lat_cent, lon_cent)
Step 4: Filter & Sum Relevant Incomes
Create a boolean mask for distances ≤ 5000 meters, then sum the corresponding income values:
# Mask where distance is <= 5000 meters distance_mask = distances <= 5000 # Sum all incomes where the mask is True total_qualified_income = np.sum(cent_incomes[distance_mask]) # If you want the sum per df row instead of a global total, use: # per_row_income_sum = np.sum(cent_incomes * distance_mask, axis=1)
Why This Is Faster
- Vectorization: NumPy operations run in optimized C code instead of slow Python loops, cutting down execution time by orders of magnitude.
- Broadcasting: Avoids creating intermediate copies of data, keeping memory usage efficient (compared to explicit nested loops).
Alternative: Using SciPy's cdist
If you prefer a more concise approach, you can use scipy.spatial.distance.cdist to compute pairwise distances with a custom metric:
from scipy.spatial.distance import cdist def haversine_metric(coord1, coord2): lat1, lon1 = coord1 lat2, lon2 = coord2 return haversine(lat1, lon1, lat2, lon2) # Compute distance matrix distances = cdist(df_coords, cent_coords, metric=haversine_metric) # Same filtering and summation as before distance_mask = distances <= 5000 total_qualified_income = np.sum(cent_incomes[distance_mask])
内容的提问来源于stack exchange,提问作者Erick Farias

