从零完成电商百万级订单可落地大数据筛选:Pandas+PySpark全步骤
时间:2026年05月31日 16:27:20
来源:易频IT社区
一、前置准备(零门槛必装组件)
1.1 本地基础环境
- 安装Python 3.9:访问https://www.python.org/downloads/release/python-3913/,选择对应系统的安装包,勾选「Add Python 3.9 to PATH」后点击「Install Now」
- 验证Python+pip:打开终端(Windows用cmd/powershell,Mac/Linux用Terminal),依次输入`python --version`、`pip --version`,返回版本号即成功
- 安装本地数据处理库:终端输入`pip install pandas numpy openpyxl -i https://pypi.tuna.tsinghua.edu.cn/simple`
1.2 本地模拟大数据环境
- 安装Java 8:访问https://www.oracle.com/java/technologies/javase/javase8-archive-downloads.html,下载「jdk-8u391-windows-x64.exe」(需注册Oracle账号,嫌麻烦可用OpenJDK 8:https://adoptium.net/zh-CN/temurin/releases/?version=8),安装后配置环境变量:
1. 新建系统变量`JAVA_HOME`,值为JDK安装路径(如`C:\Program Files\Java\jdk1.8.0_391`)
2. 编辑系统变量`Path`,新增`%JAVA_HOME%\bin`和`%JAVA_HOME%\jre\bin`
3. 终端输入`java -version`,返回1.8版本即成功
- 安装PySpark 3.3:终端输入`pip install pyspark==3.3.0 -i https://pypi.tuna.tsinghua.edu.cn/simple`
- 生成模拟百万级订单数据:新建`generate_data.py`,复制以下代码并运行(生成约1.2GB的CSV文件,3-5分钟完成)
```python
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
配置参数
NUM_ORDERS = 1_000_000
START_DATE = datetime(2024, 1, 1)
PRODUCT_CATEGORIES = ["3C数码", "服装鞋帽", "美妆个护", "食品生鲜", "家居家纺"]
PAYMENT_METHODS = ["支付宝", "微信支付", "银行卡", "花呗"]
REGIONS = ["北京", "上海", "广州", "深圳", "杭州", "成都", "武汉", "西安"]
生成随机数据
np.random.seed(42) 固定随机数,便于复现
order_ids = np.arange(1, NUM_ORDERS + 1)
user_ids = np.random.randint(1, 100_001, NUM_ORDERS)
product_ids = np.random.randint(1, 5001, NUM_ORDERS)
categories = np.random.choice(PRODUCT_CATEGORIES, NUM_ORDERS)
amounts = np.round(np.random.uniform(9.9, 9999.9, NUM_ORDERS), 2)
payment_methods = np.random.choice(PAYMENT_METHODS, NUM_ORDERS)
regions = np.random.choice(REGIONS, NUM_ORDERS)
statuses = np.random.choice(["待支付", "已支付", "已发货", "已收货", "已退款"], NUM_ORDERS, p=[0.05, 0.15, 0.3, 0.4, 0.1])
order_dates = [START_DATE + timedelta(days=np.random.randint(0, 365)) for _ in range(NUM_ORDERS)]
构建DataFrame并导出
df = pd.DataFrame({
"order_id": order_ids,
"user_id": user_ids,
"product_id": product_ids,
"category": categories,
"amount": amounts,
"payment_method": payment_methods,
"region": regions,
"status": statuses,
"order_date": order_dates
})
df.to_csv("orders_1m.csv", index=False, encoding="utf-8-sig")
print("百万级订单数据生成成功!")
```
二、本地Pandas小批量预筛选(验证筛选逻辑)
2.1 明确筛选需求
以「2024年6-8月杭州地区、已收货的美妆个护订单,金额≥500元」为例,本地读取前1000行验证逻辑
2.2 预筛选代码
新建`pandas_preview.py`,复制以下代码并运行
```python
import pandas as pd
读取前1000行小批量数据
df_small = pd.read_csv("orders_1m.csv", nrows=1000, encoding="utf-8-sig")
转换日期格式
df_small["order_date"] = pd.to_datetime(df_small["order_date"])
执行筛选逻辑
filter_df = df_small[
(df_small["order_date"].dt.month.isin([6,7,8])) & 6-8月
(df_small["region"] == "杭州") & 杭州地区
(df_small["status"] == "已收货") & 已收货
(df_small["category"] == "美妆个护") & 美妆个护
(df_small["amount"] >= 500) 金额≥500
]
导出验证结果
filter_df.to_csv("pandas_filtered_small.csv", index=False, encoding="utf-8-sig")
print(f"小批量预筛选完成,共找到{len(filter_df)}条符合条件的订单")
```
三、PySpark本地/单机集群全量筛选(处理百万级数据)
3.1 配置SparkSession
新建`pyspark_full.py`,首先导入库并创建SparkSession(单机无需额外配置集群)
```python
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, month
初始化SparkSession
spark = SparkSession.builder \
.appName("OrdersFilter") \
.master("local[]") []表示使用所有CPU核心
.getOrCreate()
关闭不必要的日志
spark.sparkContext.setLogLevel("WARN")
```
3.2 全量读取+筛选+导出
在上述代码下方继续添加:
```python
全量读取百万级CSV数据
df_full = spark.read.csv(
"orders_1m.csv",
header=True, 第一行为表头
inferSchema=True, 自动推断数据类型
encoding="utf-8-sig"
)
转换日期格式(inferSchema可能识别不完整,强制转换更稳妥)
df_full = df_full.withColumn("order_date", col("order_date").cast("date"))
执行与预筛选完全一致的逻辑
filtered_df = df_full.filter(
(month(col("order_date")).isin([6,7,8])) &
(col("region") == "杭州") &
(col("status") == "已收货") &
(col("category") == "美妆个护") &
(col("amount") >= 500)
)
查看筛选结果的前10行(可选,用于确认)
filtered_df.show(10, truncate=False)
统计符合条件的订单总数(可选)
print(f"全量筛选完成,共找到{filtered_df.count()}条符合条件的订单")
导出筛选结果(生成一个文件夹,内含多个小CSV文件,可用Excel合并工具或Pandas合并)
filtered_df.coalesce(1).write.csv(
"pyspark_filtered_full",
header=True,
encoding="utf-8-sig",
mode="overwrite" 覆盖已有文件夹
)
关闭SparkSession
spark.stop()
```
3.3 合并导出的小CSV文件
PySpark默认按分区导出,这里用`coalesce(1)`合并为1个分区,但文件名会随机生成,新建`merge_csv.py`批量读取并合并:
```python
import os
import pandas as pd
找到导出文件夹中的所有CSV文件
folder_path = "pyspark_filtered_full"
csv_files = [f for f in os.listdir(folder_path) if f.endswith(".csv")]
读取并合并
df_merged = pd.concat([pd.read_csv(os.path.join(folder_path, f), encoding="utf-8-sig") for f in csv_files])
导出最终文件
df_merged.to_csv("final_filtered_orders.csv", index=False, encoding="utf-8-sig")
print("最终合并文件生成成功!")
```
四、快速验证全量筛选结果
打开`final_filtered_orders.csv`,随机抽查3-5条数据,检查:
1. 订单日期是否在2024年6-8月
2. 地区是否为杭州
3. 状态是否为已收货
4. 品类是否为美妆个护
5. 金额是否≥500元
五、常见问题排查
1. Java版本不兼容:PySpark 3.3仅支持Java 8/11,优先使用Java 8
2. 中文乱码:所有读写操作均指定`encoding="utf-8-sig"`
3. 内存不足:若本地内存<16GB,可将SparkSession的`master`改为`local[2]`,或调整配置:
```python
spark = SparkSession.builder \
.appName("OrdersFilter") \
.master("local[2]") \
.config("spark.driver.memory", "4g") \
.config("spark.executor.memory", "4g") \
.getOrCreate()
```