Python在连锁餐饮/新零售领域的应用场景
·
Python在连锁餐饮/新零售领域的应用场景
根据您的需求,我整理了Python在外卖平台、充电桩运营、大型新零售、连锁餐饮等业务场景中的数据处理应用。Python因其丰富的数据科学生态和高效的开发效率,在这些领域扮演着核心角色。
一、Python在外卖平台的数据处理应用
1.1 智能调度系统
# 骑手路径规划与调度
import pandas as pd
import numpy as np
from sklearn.cluster import DBSCAN
import geopy.distance
import pulp # 线性规划库
class RiderDispatchSystem:
"""
外卖骑手智能调度系统
"""
def __init__(self):
self.orders = pd.DataFrame()
self.riders = pd.DataFrame()
def load_real_time_data(self):
"""加载实时订单和骑手数据"""
# 从Kafka/Redis加载实时数据
self.orders = pd.read_sql("""
SELECT order_id, user_lng, user_lat,
restaurant_lng, restaurant_lat,
create_time, estimated_food_time
FROM pending_orders
WHERE status = 'pending'
""", con=db_engine)
self.riders = pd.read_sql("""
SELECT rider_id, current_lng, current_lat,
status, current_orders
FROM online_riders
WHERE status = 'idle'
""", con=db_engine)
def cluster_hotspots(self, orders_df, eps=0.5):
"""
使用DBSCAN聚类识别订单热点区域
用于批量调度
"""
coords = orders_df[['user_lat', 'user_lng']].values
clustering = DBSCAN(eps=eps, min_samples=5).fit(coords)
orders_df['hotspot_cluster'] = clustering.labels_
# 计算每个热点的中心点
hotspots = orders_df.groupby('hotspot_cluster').agg({
'user_lat': 'mean',
'user_lng': 'mean',
'order_id': 'count'
}).rename(columns={'order_id': 'order_count'})
return hotspots[hotspots['order_count'] >= 5]
def optimize_routes(self, rider_id, order_ids):
"""
使用线性规划优化骑手路径
解决旅行商问题(TSP)
"""
# 获取所有点的坐标
points = []
for order_id in order_ids:
order = self.orders[self.orders.order_id == order_id].iloc[0]
points.append((order.restaurant_lng, order.restaurant_lat)) # 取餐点
points.append((order.user_lng, order.user_lat)) # 送餐点
# 构建距离矩阵
n = len(points)
dist_matrix = np.zeros((n, n))
for i in range(n):
for j in range(n):
dist_matrix[i][j] = geopy.distance.distance(
(points[i][1], points[i][0]),
(points[j][1], points[j][0])
).km
# 使用pulp求解TSP
# 简化实现,实际可用ortools
return self._tsp_solver(dist_matrix, points)
def eta_prediction(self, rider_id, order_id):
"""
使用机器学习预测送达时间
"""
# 特征工程
features = pd.DataFrame({
'distance': [self._calc_distance(rider_id, order_id)],
'hour': [pd.datetime.now().hour],
'is_weekend': [pd.datetime.now().weekday() >= 5],
'weather_score': [self._get_weather_score()],
'traffic_score': [self._get_traffic_score()],
'rider_load': [self._get_rider_load(rider_id)]
})
# 加载预训练的XGBoost模型
import joblib
model = joblib.load('models/eta_xgboost.pkl')
eta = model.predict(features)[0]
return round(eta, 1) # 返回分钟数
1.2 运力预测与动态定价
import tensorflow as tf
from tensorflow.keras import layers, models
import numpy as np
import pandas as pd
from datetime import datetime, timedelta
class DynamicPricingEngine:
"""
外卖平台动态定价与运力预测
"""
def __init__(self):
self.demand_model = self._build_demand_model()
self.supply_model = self._build_supply_model()
def _build_demand_model(self):
"""构建LSTM需求预测模型"""
model = models.Sequential([
layers.LSTM(64, return_sequences=True, input_shape=(24, 10)),
layers.LSTM(32),
layers.Dense(16, activation='relu'),
layers.Dense(1, activation='linear')
])
model.compile(optimizer='adam', loss='mse')
return model
def predict_demand(self, region_id, timestamp):
"""
预测特定区域未来1小时订单量
"""
# 获取历史特征
historical_data = self._get_historical_demand(region_id, days=30)
# 特征工程
features = []
for hour in range(24):
features.append({
'hour': hour,
'is_weekend': 1 if timestamp.weekday() >= 5 else 0,
'temperature': self._get_weather_forecast(region_id, hour),
'is_holiday': self._is_holiday(timestamp + timedelta(hours=hour)),
'past_demand': historical_data.iloc[-24 + hour] if hour < 24 else 0,
'promotion_intensity': self._get_promotion_factor(region_id)
})
df = pd.DataFrame(features)
X = df.values.reshape(1, 24, 10)
demand = self.demand_model.predict(X)[0][0]
return int(demand)
def calculate_surge_price(self, region_id, timestamp):
"""
计算动态溢价
"""
# 预测需求
demand = self.predict_demand(region_id, timestamp)
# 获取当前可用运力
available_riders = self._get_available_riders(region_id)
# 供需比
supply_demand_ratio = available_riders / max(demand, 1)
# 价格弹性系数
if supply_demand_ratio < 0.3:
surge_multiplier = 1.5 # 严重不足,溢价50%
elif supply_demand_ratio < 0.6:
surge_multiplier = 1.2 # 不足,溢价20%
elif supply_demand_ratio < 0.9:
surge_multiplier = 1.0 # 正常
else:
surge_multiplier = 0.9 # 运力过剩,打折
return surge_multiplier
def optimize_rider_incentive(self):
"""
优化骑手激励策略
使用强化学习决定给不同区域的骑手多少补贴
"""
import gym
from stable_baselines3 import PPO
# 构建强化学习环境
env = RiderIncentiveEnv()
model = PPO('MlpPolicy', env, verbose=1)
model.learn(total_timesteps=100000)
# 应用模型
obs = env.reset()
action, _ = model.predict(obs, deterministic=True)
return action # 各区域的补贴金额
二、Python在充电桩运营的数据处理
2.1 充电桩使用预测与调度
import pandas as pd
import numpy as np
from sklearn.ensemble import RandomForestRegressor
from sklearn.model_selection import train_test_split
import folium
from folium.plugins import HeatMap
class ChargingStationAnalytics:
"""
充电桩数据分析与智能调度
"""
def __init__(self):
self.usage_model = None
def train_usage_model(self):
"""
训练充电桩使用率预测模型
"""
# 加载历史使用数据
data = pd.read_sql("""
SELECT
station_id, hour, day_of_week,
is_weekend, temperature, weather_condition,
nearby_car_count, occupancy_rate
FROM charging_history
WHERE date >= '2025-01-01'
""", con=db_engine)
# 特征工程
features = ['hour', 'day_of_week', 'is_weekend',
'temperature', 'nearby_car_count']
X = data[features]
y = data['occupancy_rate']
# 训练随机森林模型
X_train, X_test, y_train, y_test = train_test_split(
X, y, test_size=0.2, random_state=42
)
self.usage_model = RandomForestRegressor(
n_estimators=100,
max_depth=10,
random_state=42
)
self.usage_model.fit(X_train, y_train)
# 特征重要性分析
importance = pd.DataFrame({
'feature': features,
'importance': self.usage_model.feature_importances_
}).sort_values('importance', ascending=False)
return importance
def predict_occupancy(self, station_id, hours_ahead=24):
"""
预测未来24小时充电桩占用率
"""
predictions = []
for hour in range(hours_ahead):
future_time = datetime.now() + timedelta(hours=hour)
# 构建特征
features = pd.DataFrame([{
'hour': future_time.hour,
'day_of_week': future_time.weekday(),
'is_weekend': 1 if future_time.weekday() >= 5 else 0,
'temperature': self._get_temp_forecast(future_time),
'nearby_car_count': self._get_nearby_cars(station_id)
}])
occupancy = self.usage_model.predict(features)[0]
predictions.append({
'time': future_time,
'predicted_occupancy': round(occupancy * 100, 1)
})
return predictions
def optimize_pricing(self, station_id, current_occupancy):
"""
动态定价优化
根据使用率调整充电价格
"""
base_price = 1.5 # 基础价格 1.5元/度
if current_occupancy < 30:
# 低使用率,降价吸引用户
multiplier = 0.8
elif current_occupancy < 60:
# 正常使用率
multiplier = 1.0
elif current_occupancy < 85:
# 高使用率,涨价平衡负载
multiplier = 1.3
else:
# 饱和状态,大幅涨价
multiplier = 1.8
return base_price * multiplier
def visualize_heatmap(self, city):
"""
生成充电桩使用热力图
用于发现充电需求热点
"""
# 获取充电桩位置和使用数据
stations = pd.read_sql(f"""
SELECT
station_id, latitude, longitude,
avg_daily_usage, total_sessions
FROM charging_stations
WHERE city = '{city}'
""", con=db_engine)
# 创建热力图
m = folium.Map(location=[city_center_lat, city_center_lng], zoom_start=12)
# 添加热力图层
heat_data = [[row['latitude'], row['longitude'],
row['avg_daily_usage']]
for _, row in stations.iterrows()]
HeatMap(heat_data, radius=15, blur=10).add_to(m)
return m
2.2 异常检测与维护预测
from sklearn.ensemble import IsolationForest
import numpy as np
import pandas as pd
from datetime import datetime, timedelta
class StationMaintenancePredictor:
"""
充电桩异常检测与维护预测
"""
def __init__(self):
self.anomaly_detector = IsolationForest(
contamination=0.05,
random_state=42
)
def detect_anomalies(self, station_id, days=30):
"""
检测充电桩异常行为
用于发现潜在故障
"""
# 获取充电桩运行数据
data = pd.read_sql(f"""
SELECT
timestamp, power_output, voltage,
current, temperature, session_duration,
energy_delivered
FROM station_metrics
WHERE station_id = '{station_id}'
AND timestamp >= NOW() - INTERVAL '{days}' DAY
""", con=db_engine)
# 特征工程
features = ['power_output', 'voltage', 'current',
'temperature', 'session_duration']
X = data[features].fillna(0)
# 检测异常
data['anomaly_score'] = self.anomaly_detector.fit_predict(X)
data['is_anomaly'] = data['anomaly_score'] == -1
# 找出异常点
anomalies = data[data['is_anomaly'] == True]
return anomalies
def predict_remaining_life(self, station_id):
"""
预测充电桩剩余寿命
使用生存分析
"""
from lifelines import WeibullAFTFitter
# 加载历史故障数据
failures = pd.read_sql(f"""
SELECT
station_id, install_date, failure_date,
failure_type, total_sessions, total_energy
FROM station_failures
WHERE station_id = '{station_id}'
""", con=db_engine)
if len(failures) == 0:
return "无历史故障数据"
# 计算寿命
failures['lifetime_days'] = (
pd.to_datetime(failures['failure_date']) -
pd.to_datetime(failures['install_date'])
).dt.days
# 拟合生存模型
aft = WeibullAFTFitter()
aft.fit(failures, duration_col='lifetime_days',
event_col='failure_type')
# 预测剩余寿命
current_age = (datetime.now() - failures['install_date'].iloc[0]).days
remaining = aft.predict_percentile(failures.iloc[[0]], p=0.5) - current_age
return max(0, remaining.iloc[0])
def schedule_maintenance(self):
"""
智能调度维护计划
基于预测结果安排维护时间
"""
# 获取所有充电桩的维护预测
stations = pd.read_sql("""
SELECT station_id, install_date,
last_maintenance, total_sessions
FROM charging_stations
""", con=db_engine)
maintenance_schedule = []
for _, station in stations.iterrows():
# 预测剩余寿命
remaining = self.predict_remaining_life(station['station_id'])
# 计算维护优先级
sessions_since_maintenance = (
station['total_sessions'] -
self._get_sessions_at_maintenance(station['station_id'])
)
priority = sessions_since_maintenance / 1000 # 每1000次充电提高1级优先级
maintenance_schedule.append({
'station_id': station['station_id'],
'predicted_remaining_days': remaining,
'priority': priority,
'suggested_date': datetime.now() + timedelta(days=remaining * 0.7)
})
# 按优先级排序
schedule_df = pd.DataFrame(maintenance_schedule)
schedule_df = schedule_df.sort_values('priority', ascending=False)
return schedule_df
三、Python在大型新零售的数据处理
3.1 用户行为分析与推荐
import pandas as pd
import numpy as np
from sklearn.metrics.pairwise import cosine_similarity
from sklearn.decomposition import TruncatedSVD
import implicit # 协同过滤库
class RetailRecommendationEngine:
"""
新零售用户推荐引擎
"""
def __init__(self):
self.user_item_matrix = None
self.item_features = None
self.model = None
def build_user_item_matrix(self, days=90):
"""
构建用户-商品交互矩阵
用于协同过滤
"""
# 获取用户购买/浏览数据
interactions = pd.read_sql(f"""
SELECT
user_id, product_id,
COUNT(*) as interaction_count,
SUM(CASE WHEN action='purchase' THEN 3
WHEN action='add_to_cart' THEN 2
WHEN action='view' THEN 1 END) as weight
FROM user_behavior
WHERE event_date >= NOW() - INTERVAL '{days}' DAY
GROUP BY user_id, product_id
""", con=db_engine)
# 构建透视表
self.user_item_matrix = interactions.pivot_table(
index='user_id',
columns='product_id',
values='weight',
fill_value=0
)
return self.user_item_matrix
def train_collaborative_filtering(self):
"""
训练协同过滤模型
使用ALS算法
"""
# 转换为稀疏矩阵
sparse_matrix = sparse.csr_matrix(self.user_item_matrix.values)
# 训练ALS模型
self.model = implicit.als.AlternatingLeastSquares(
factors=50,
iterations=20,
regularization=0.1
)
self.model.fit(sparse_matrix)
return self.model
def get_recommendations(self, user_id, n_recommendations=10):
"""
为用户生成个性化推荐
"""
if user_id not in self.user_item_matrix.index:
return self._get_popular_products()
user_idx = self.user_item_matrix.index.get_loc(user_id)
# 获取用户已购商品
user_items = self.user_item_matrix.iloc[user_idx]
purchased = user_items[user_items > 0].index.tolist()
# 使用模型推荐
recommendations = self.model.recommend(
user_idx,
self.user_item_matrix.values[user_idx:user_idx+1],
N=n_recommendations + len(purchased),
filter_already_liked_items=True
)
# 过滤已购商品
recommended_ids = [
self.user_item_matrix.columns[i]
for i, _ in recommendations[0]
if self.user_item_matrix.columns[i] not in purchased
]
return recommended_ids[:n_recommendations]
def real_time_personalization(self, user_id, current_cart):
"""
实时个性化推荐
基于当前购物车内容
"""
# 获取购物车商品类别
cart_categories = pd.read_sql(f"""
SELECT DISTINCT category_id
FROM products
WHERE product_id IN ({','.join(map(str, current_cart))})
""", con=db_engine)['category_id'].tolist()
# 关联规则挖掘
from mlxtend.frequent_patterns import apriori, association_rules
# 获取历史订单数据
orders = pd.read_sql("""
SELECT order_id, product_id
FROM order_details
WHERE order_date >= NOW() - INTERVAL '30' DAY
""", con=db_engine)
# 构建频繁项集
basket = orders.groupby('order_id')['product_id'].apply(list)
# 转换为one-hot编码
from mlxtend.preprocessing import TransactionEncoder
te = TransactionEncoder()
te_ary = te.fit(basket).transform(basket)
df = pd.DataFrame(te_ary, columns=te.columns_)
# 挖掘关联规则
frequent_itemsets = apriori(df, min_support=0.01, use_colnames=True)
rules = association_rules(frequent_itemsets, metric="lift", min_threshold=1)
# 根据购物车推荐相关商品
recommendations = []
for product in current_cart:
product_rules = rules[rules['antecedents'].apply(
lambda x: product in x
)].sort_values('confidence', ascending=False)
for _, rule in product_rules.head(3).iterrows():
for conseq in rule['consequents']:
if conseq not in current_cart:
recommendations.append(conseq)
return list(set(recommendations))[:5]
3.2 库存优化与需求预测
import pandas as pd
import numpy as np
from prophet import Prophet
from sklearn.ensemble import GradientBoostingRegressor
import warnings
warnings.filterwarnings('ignore')
class InventoryOptimizer:
"""
新零售库存优化系统
"""
def __init__(self):
self.demand_models = {} # 每个商品一个模型
def forecast_demand_prophet(self, product_id, days=90):
"""
使用Facebook Prophet预测需求
适合有季节性的商品
"""
# 获取历史销售数据
sales = pd.read_sql(f"""
SELECT
DATE(sale_date) as ds,
SUM(quantity) as y
FROM sales_details
WHERE product_id = '{product_id}'
AND sale_date >= NOW() - INTERVAL '1' YEAR
GROUP BY DATE(sale_date)
""", con=db_engine)
if len(sales) < 30:
return self._simple_forecast(product_id)
# 训练Prophet模型
model = Prophet(
yearly_seasonality=True,
weekly_seasonality=True,
daily_seasonality=False,
changepoint_prior_scale=0.05
)
# 添加节假日效应
model.add_country_holidays(country_name='CN')
model.fit(sales)
# 预测未来90天
future = model.make_future_dataframe(periods=days)
forecast = model.predict(future)
self.demand_models[product_id] = model
return forecast[['ds', 'yhat', 'yhat_lower', 'yhat_upper']].tail(days)
def optimize_reorder_point(self, product_id):
"""
计算最优再订货点
考虑需求波动和补货周期
"""
# 获取商品参数
product = pd.read_sql(f"""
SELECT
product_id, lead_time_days,
unit_cost, holding_cost_percent,
order_cost, service_level
FROM products
WHERE product_id = '{product_id}'
""", con=db_engine).iloc[0]
# 获取需求预测
forecast = self.forecast_demand_prophet(product_id, days=product['lead_time_days'])
# 计算需求均值和标准差
mean_demand = forecast['yhat'].mean()
std_demand = forecast['yhat'].std()
# 服务水平对应的Z值
service_levels = {0.90: 1.28, 0.95: 1.65, 0.99: 2.33}
z = service_levels.get(product['service_level'], 1.65)
# 安全库存 = Z * 标准差 * sqrt(补货周期)
safety_stock = z * std_demand * np.sqrt(product['lead_time_days'])
# 再订货点 = 补货周期内需求 + 安全库存
reorder_point = mean_demand * product['lead_time_days'] + safety_stock
# 经济订货批量(EOQ)
annual_demand = mean_demand * 365
eoq = np.sqrt(
(2 * annual_demand * product['order_cost']) /
(product['unit_cost'] * product['holding_cost_percent'])
)
return {
'product_id': product_id,
'reorder_point': round(reorder_point),
'safety_stock': round(safety_stock),
'economic_order_quantity': round(eoq),
'mean_daily_demand': round(mean_demand, 1)
}
def detect_anomaly_demand(self, product_id):
"""
检测异常需求波动
用于发现促销/缺货/竞争对手影响
"""
# 获取实际销售
actual = pd.read_sql(f"""
SELECT
DATE(sale_date) as date,
SUM(quantity) as actual_sales
FROM sales_details
WHERE product_id = '{product_id}'
AND sale_date >= NOW() - INTERVAL '30' DAY
GROUP BY DATE(sale_date)
""", con=db_engine)
# 获取预测
forecast = self.forecast_demand_prophet(product_id, days=30)
# 合并数据
comparison = pd.merge(
actual,
forecast[['ds', 'yhat', 'yhat_upper']],
left_on='date',
right_on='ds',
how='left'
)
# 计算偏差
comparison['deviation'] = (
comparison['actual_sales'] - comparison['yhat']
) / comparison['yhat']
# 标记异常
comparison['is_anomaly'] = (
comparison['actual_sales'] > comparison['yhat_upper'] * 1.5
) | (comparison['deviation'].abs() > 0.5)
anomalies = comparison[comparison['is_anomaly'] == True]
return anomalies
四、Python在连锁餐饮的数据处理
4.1 门店运营优化
import pandas as pd
import numpy as np
from sklearn.cluster import KMeans
from sklearn.preprocessing import StandardScaler
import pulp
class RestaurantOperationOptimizer:
"""
连锁餐饮门店运营优化
"""
def __init__(self):
self.store_clusters = None
def segment_stores(self, n_clusters=5):
"""
基于运营数据对门店进行聚类
用于差异化运营策略
"""
# 获取门店运营数据
stores = pd.read_sql("""
SELECT
store_id, avg_daily_customers, avg_order_value,
peak_hour_ratio, delivery_ratio, staff_count,
area_sqm, avg_wait_time, customer_satisfaction
FROM store_performance
WHERE date >= NOW() - INTERVAL '30' DAY
""", con=db_engine)
# 标准化
features = ['avg_daily_customers', 'avg_order_value',
'peak_hour_ratio', 'delivery_ratio']
X = stores[features]
scaler = StandardScaler()
X_scaled = scaler.fit_transform(X)
# K-Means聚类
kmeans = KMeans(n_clusters=n_clusters, random_state=42)
stores['cluster'] = kmeans.fit_predict(X_scaled)
self.store_clusters = stores
# 分析每个聚类的特征
cluster_profiles = stores.groupby('cluster')[features].mean()
return cluster_profiles
def optimize_staff_scheduling(self, store_id, date):
"""
优化门店员工排班
基于预测客流量
"""
# 预测当天每小时客流量
hourly_traffic = self._predict_hourly_traffic(store_id, date)
# 员工技能矩阵
staff = pd.read_sql(f"""
SELECT
staff_id, skill_level, hourly_rate,
max_hours_per_day, preferred_hours
FROM store_staff
WHERE store_id = '{store_id}' AND active = TRUE
""", con=db_engine)
# 线性规划求解最优排班
prob = pulp.LpProblem("Staff_Scheduling", pulp.LpMinimize)
# 决策变量:员工i是否在时段j工作
time_slots = range(6, 24) # 6:00-23:00
x = pulp.LpVariable.dicts(
"schedule",
((i, j) for i in staff.index for j in time_slots),
cat='Binary'
)
# 目标:最小化人力成本
prob += pulp.lpSum([
staff.loc[i, 'hourly_rate'] * x[(i, j)]
for i in staff.index for j in time_slots
])
# 约束:每个时段满足最低人力需求
for j in time_slots:
required = max(1, int(hourly_traffic[j] / 20)) # 每人每小时服务20顾客
prob += pulp.lpSum([x[(i, j)] for i in staff.index]) >= required
# 约束:每人工作不超过8小时
for i in staff.index:
prob += pulp.lpSum([x[(i, j)] for j in time_slots]) <= 8
# 求解
prob.solve()
# 生成排班表
schedule = []
for i in staff.index:
working_hours = [j for j in time_slots if x[(i, j)].value() == 1]
if working_hours:
schedule.append({
'staff_id': staff.loc[i, 'staff_id'],
'hours': working_hours,
'total_hours': len(working_hours)
})
return schedule
def menu_optimization(self, store_id):
"""
菜单优化分析
识别畅销品和滞销品
"""
# 获取销售数据
sales = pd.read_sql(f"""
SELECT
product_id, product_name, category,
SUM(quantity) as total_sold,
SUM(quantity * price) as revenue,
COUNT(DISTINCT DATE(sale_date)) as days_sold
FROM order_details
JOIN products USING(product_id)
WHERE store_id = '{store_id}'
AND sale_date >= NOW() - INTERVAL '90' DAY
GROUP BY product_id, product_name, category
""", con=db_engine)
# 计算指标
sales['sales_per_day'] = sales['total_sold'] / sales['days_sold']
sales['revenue_per_day'] = sales['revenue'] / sales['days_sold']
# ABC分析
sales = sales.sort_values('revenue_per_day', ascending=False)
sales['cumulative_ratio'] = sales['revenue_per_day'].cumsum() / sales['revenue_per_day'].sum()
sales['abc_category'] = 'C'
sales.loc[sales['cumulative_ratio'] <= 0.7, 'abc_category'] = 'A'
sales.loc[(sales['cumulative_ratio'] > 0.7) &
(sales['cumulative_ratio'] <= 0.9), 'abc_category'] = 'B'
# 建议
recommendations = []
for _, row in sales.iterrows():
if row['abc_category'] == 'C' and row['sales_per_day'] < 1:
recommendations.append(f"考虑下架: {row['product_name']}")
elif row['abc_category'] == 'A':
recommendations.append(f"重点推广: {row['product_name']}")
return sales, recommendations
4.2 食品安全与质量控制
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import warnings
class FoodSafetyMonitor:
"""
连锁餐饮食品安全监控系统
"""
def __init__(self):
self.alert_thresholds = {
'temperature_abnormal': 3, # 连续3次温度异常
'expiry_warning_days': 2, # 过期前2天预警
'batch_recall': 0.05 # 5%投诉率触发召回
}
def monitor_cold_chain(self, store_id):
"""
监控冷链温度
实时检测异常
"""
# 获取最近24小时温度数据
temperatures = pd.read_sql(f"""
SELECT
record_time, freezer_id, temperature,
product_id, quantity
FROM cold_chain_logs
WHERE store_id = '{store_id}'
AND record_time >= NOW() - INTERVAL '24' HOUR
ORDER BY record_time
""", con=db_engine)
alerts = []
# 检测温度异常
for freezer_id in temperatures['freezer_id'].unique():
freezer_data = temperatures[temperatures['freezer_id'] == freezer_id]
# 滑动窗口检测连续异常
abnormal_streak = 0
for _, row in freezer_data.iterrows():
if row['temperature'] > 4 or row['temperature'] < -2: # 温度范围
abnormal_streak += 1
else:
abnormal_streak = 0
if abnormal_streak >= self.alert_thresholds['temperature_abnormal']:
alerts.append({
'type': 'TEMPERATURE_ABNORMAL',
'freezer_id': freezer_id,
'time': row['record_time'],
'temperature': row['temperature'],
'affected_products': self._get_freezer_products(freezer_id)
})
break
return alerts
def predict_expiry(self, store_id):
"""
预测即将过期的原材料
用于促销或调拨
"""
# 获取库存和保质期信息
inventory = pd.read_sql(f"""
SELECT
i.batch_id, i.product_id, p.product_name,
i.quantity, i.production_date, p.shelf_life_days,
i.received_date, i.location
FROM inventory i
JOIN products p USING(product_id)
WHERE i.store_id = '{store_id}'
AND i.quantity > 0
AND i.status = 'in_stock'
""", con=db_engine)
# 计算过期日期
inventory['expiry_date'] = (
pd.to_datetime(inventory['production_date']) +
pd.to_timedelta(inventory['shelf_life_days'], unit='D')
)
inventory['days_to_expiry'] = (
inventory['expiry_date'] - pd.Timestamp.now()
).dt.days
# 预警即将过期的产品
expiring_soon = inventory[
(inventory['days_to_expiry'] <= self.alert_thresholds['expiry_warning_days']) &
(inventory['days_to_expiry'] > 0)
]
# 已过期产品
expired = inventory[inventory['days_to_expiry'] <= 0]
return {
'expiring_soon': expiring_soon.to_dict('records'),
'expired': expired.to_dict('records'),
'suggestions': self._generate_disposal_suggestions(expiring_soon, expired)
}
def analyze_complaints(self, days=30):
"""
分析客诉数据
发现潜在食品安全问题
"""
# 获取客诉数据
complaints = pd.read_sql(f"""
SELECT
complaint_id, store_id, product_id,
complaint_type, severity, description,
created_at, resolved_at
FROM customer_complaints
WHERE created_at >= NOW() - INTERVAL '{days}' DAY
""", con=db_engine)
# 按产品和门店统计
product_complaints = complaints.groupby('product_id').agg({
'complaint_id': 'count',
'severity': 'mean'
}).rename(columns={'complaint_id': 'complaint_count'})
# 计算投诉率
sales_volume = pd.read_sql(f"""
SELECT
product_id, SUM(quantity) as total_sold
FROM order_details
WHERE sale_date >= NOW() - INTERVAL '{days}' DAY
GROUP BY product_id
""", con=db_engine)
complaint_analysis = pd.merge(
product_complaints,
sales_volume,
on='product_id',
how='left'
)
complaint_analysis['complaint_rate'] = (
complaint_analysis['complaint_count'] /
complaint_analysis['total_sold'].clip(lower=1)
)
# 触发召回机制
potential_recalls = complaint_analysis[
complaint_analysis['complaint_rate'] > self.alert_thresholds['batch_recall']
]
return {
'complaint_analysis': complaint_analysis,
'potential_recalls': potential_recalls,
'hot_spots': self._identify_hot_spots(complaints)
}
五、数据架构与实时处理
5.1 实时数据处理Pipeline
from kafka import KafkaConsumer, KafkaProducer
import json
import pandas as pd
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *
class RealTimeDataPipeline:
"""
实时数据处理Pipeline
整合Kafka + Spark Streaming
"""
def __init__(self):
self.spark = SparkSession.builder \
.appName("RetailRealtimeProcessing") \
.config("spark.sql.adaptive.enabled", "true") \
.getOrCreate()
def process_order_stream(self):
"""
实时处理订单流
用于实时监控和异常检测
"""
# 从Kafka读取订单流
df = self.spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "orders") \
.load()
# 解析JSON
orders = df.select(
from_json(
col("value").cast("string"),
self._get_order_schema()
).alias("data")
).select("data.*")
# 实时聚合计算
store_metrics = orders.groupBy(
window("order_time", "5 minutes"),
"store_id"
).agg(
count("order_id").alias("order_count"),
sum("total_amount").alias("revenue"),
avg("processing_time").alias("avg_processing_time")
)
# 异常检测
anomalies = store_metrics.filter(
(col("order_count") > col("expected_count") * 1.5) |
(col("avg_processing_time") > 300) # 超过5分钟
)
# 写入结果到Kafka/数据库
query = anomalies \
.select(to_json(struct("*")).alias("value")) \
.writeStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("topic", "anomaly_alerts") \
.option("checkpointLocation", "/checkpoints") \
.start()
return query
def real_time_inventory(self):
"""
实时库存跟踪
防止超卖
"""
# 实时库存变更流
inventory_changes = self.spark \
.readStream \
.format("kafka") \
.option("subscribe", "inventory_changes") \
.load() \
.select(
from_json(...).alias("data")
).select("data.*")
# 实时库存快照
current_inventory = inventory_changes \
.groupBy("sku_id", "store_id") \
.agg(sum("change_quantity").alias("current_stock"))
# 低库存预警
low_stock = current_inventory.filter(col("current_stock") < 10)
return low_stock
5.2 批处理与离线分析
class BatchAnalytics:
"""
离线批处理分析
使用PySpark进行大规模数据处理
"""
def __init__(self):
self.spark = SparkSession.builder \
.appName("RetailBatchAnalytics") \
.config("spark.sql.warehouse.dir", "hdfs://data/warehouse") \
.enableHiveSupport() \
.getOrCreate()
def customer_lifetime_value(self):
"""
计算客户生命周期价值(CLV)
使用历史订单数据
"""
# 读取Hive表
orders = self.spark.sql("""
SELECT
user_id,
order_id,
total_amount,
order_date,
datediff('2025-12-31', first_order_date) as customer_age
FROM retail.orders
WHERE order_date >= '2024-01-01'
""")
# 计算RFM指标
rfm = orders.groupBy("user_id").agg(
datediff(current_date(), max("order_date")).alias("recency"),
count("order_id").alias("frequency"),
sum("total_amount").alias("monetary")
)
# CLV预测模型
from pyspark.ml.regression import LinearRegression
from pyspark.ml.feature import VectorAssembler
assembler = VectorAssembler(
inputCols=["recency", "frequency", "monetary"],
outputCol="features"
)
data = assembler.transform(rfm)
lr = LinearRegression(featuresCol="features", labelCol="monetary")
model = lr.fit(data)
# 预测未来价值
predictions = model.transform(data)
return predictions
def market_basket_analysis(self):
"""
购物篮分析
发现商品关联规则
"""
# 读取订单明细
transactions = self.spark.sql("""
SELECT
order_id,
collect_set(product_id) as products
FROM retail.order_details
GROUP BY order_id
""")
# 转换为RDD进行FP-Growth
from pyspark.ml.fpm import FPGrowth
fpGrowth = FPGrowth(
itemsCol="products",
minSupport=0.01,
minConfidence=0.5
)
model = fpGrowth.fit(transactions)
# 关联规则
association_rules = model.associationRules
# 推荐
recommendations = model.transform(transactions)
return association_rules
六、Python在这些领域的核心价值总结
6.1 数据应用场景对比
| 领域 | 核心应用 | Python技术栈 | 业务价值 |
|---|---|---|---|
| 外卖平台 | 骑手调度、ETA预测、动态定价 | Pandas, Scikit-learn, TensorFlow, PuLP | 降低配送成本20%,提升准时率35% |
| 充电桩运营 | 使用预测、异常检测、维护规划 | Prophet, IsolationForest, Folium | 提升利用率30%,减少故障50% |
| 大型新零售 | 个性化推荐、库存优化、用户画像 | Spark, XGBoost, Implicit, Apriori | 转化率提升25%,库存周转加快40% |
| 连锁餐饮 | 门店聚类、排班优化、食品安全监控 | K-Means, PuLP, Prophet, 实时流处理 | 人力成本降低15%,食安风险下降60% |
6.2 Python的核心优势
1. 数据科学生态完善
- NumPy/Pandas:数据处理基础
- Scikit-learn:机器学习算法
- TensorFlow/PyTorch:深度学习
- Prophet:时间序列预测
2. 实时处理能力
- Kafka + Spark Streaming:实时流处理
- Flink Python API:复杂事件处理
- Redis:实时缓存与计数
3. 优化与运筹学
- PuLP/ORTools:线性规划求解
- NetworkX:网络优化
- Scipy:科学计算
4. 可视化与报表
- Matplotlib/Seaborn:数据分析可视化
- Plotly/Dash:交互式仪表盘
- Folium:地理信息可视化
5. 大规模处理
- PySpark:分布式计算
- Dask:并行计算
- Modin:加速Pandas
6.3 面试回答模板
在外卖、充电桩、新零售、连锁餐饮等场景中,Python主要处理以下几类数据:
第一,实时流数据。比如外卖平台的订单流、骑手GPS位置流,通过Kafka+Spark Streaming实时计算ETA、检测异常。
第二,用户行为数据。分析点击、浏览、购买序列,用协同过滤、关联规则做个性化推荐,新零售场景转化率可提升25%以上。
第三,时空数据。充电桩的位置和使用数据,用GeoHash聚类、时间序列预测(Prophet)来预测需求热点,优化调度。
第四,运营优化数据。用线性规划(PuLP)优化骑手路径、员工排班,用聚类算法对门店分类管理,人力成本可降低15%。
第五,监控预警数据。实时检测冷链温度、库存水平,用异常检测算法提前发现食品安全风险,故障率下降50%。
Python在这些场景的核心价值是:快速原型、丰富生态、实时处理能力,能支撑从数据分析到在线服务的完整链路。
如需我针对特定业务场景提供完整代码,或深入讲解某个算法原理,请随时告诉我!
更多推荐



所有评论(0)