项目背景与介绍
项目背景
在数字化电商时代,用户行为分析平台作为电商系统的核心组件,承担着用户画像分析、行为预测、个性化推荐、流失预警等重要功能。本项目基于国产操作系统环境,采用现代化的大数据技术栈,构建了一个功能完整、性能优异的智能电商分析平台。
项目介绍
本项目是一个基于大数据技术栈开发的电商用户行为分析平台,采用微服务架构设计。系统具备以下核心功能:
- 实时用户行为监控:基于 Kafka 的实时数据流处理、用户行为实时展示。
- 数据分析引擎:基于 Spark 的离线数据处理、Hive 数据仓库分析。
- 智能分析:基于机器学习的用户分群、流失预测、商品推荐。
- 可视化展示:基于 ECharts 的多维度数据可视化、交互式图表。
- Web 服务:基于 FastAPI 的 RESTful API、实时数据推送。
- 前端界面:响应式设计、实时数据更新、用户友好界面。
技术栈
- 大数据平台:Hadoop HDFS + YARN + Spark + Hive + HBase + Kafka。
- 后端开发:FastAPI + Pydantic + Uvicorn + SQLite/MariaDB。
- AI 算法:PyTorch + Scikit-learn + 机器学习预测模型。
- 前端技术:HTML5 + CSS3 + JavaScript + ECharts。
- 开发环境:Python 3.12 + Anaconda + Docker + Ubuntu。
项目结构
ecommerce-user-behavior/
├── src/ # 核心源码
│ ├── data_models.py # 数据模型定义
│ ├── data_generator.py # 实时数据生成器
│ ├── data_processor.py # 数据处理器
│ └── ml_analyzer.py # 机器学习分析器
├── web/ # Web 服务
│ ├── app.py # FastAPI 主应用
│ └── static/ # 静态文件
│ └── index.html # 前端页面
├── config/ # 配置文件
│ ├── settings.py # 项目配置
│ └── config.py # 系统配置
├── data/ # 数据存储
│ ├── raw/ # 原始数据
│ └── processed/ # 处理后数据
├── models/ # 模型文件
├── scripts/ # 脚本文件
│ ├── start_project.sh # 项目启动脚本
│ └── stop_project.sh # 项目停止脚本
├── logs/ # 日志文件
├── requirements.txt # Python 依赖
└── README.md # 项目说明
实操比赛题目
模块一:大数据平台搭建(20%)
任务一:基于国产操作系统搭建大数据平台
题目 1.1
在国产操作系统环境中,需要配置 Hadoop 环境变量。请完成以下操作:
# 1. 设置 JAVA_HOME 环境变量
export JAVA_HOME=____
# 2. 设置 HADOOP_HOME 环境变量
export HADOOP_HOME=____
# 3. 设置 HDFS 配置目录
export HADOOP_CONF_DIR=____
# 4. 将 Hadoop 命令添加到 PATH
export PATH=$PATH:____/bin:____/sbin
答案:____ ____ ____ ____ ____ (每空 1 分,共 5 分)
题目 1.2
在 Hadoop 配置文件中,需要设置正确的 HDFS 配置。请在 hdfs-site.xml 中补充:
<property>
<name>dfs.replication</name>
<value>____</value>
</property>
<property>
<name>dfs.namenode.name.dir</name>
<value>____</value>
</property>
<property>
<name>dfs.datanode.data.dir</name>
<value>____</value>
</property>
<property>
<name>dfs.namenode.http-address</name>
<value>____</value>
</property>
答案:____ ____ ____ ____ (每空 1 分,共 4 分)
任务二:搭建离线处理平台
题目 2.1
在 Spark 配置中,需要设置正确的 Master URL。请在 spark-defaults.conf 中补充:
# Spark 配置
spark.master=____
spark.app.name=____
spark.driver.memory=____
spark.executor.memory=____
答案:____ ____ ____ ____ (每空 1 分,共 4 分)
题目 2.2
在 Hive 配置中,需要设置正确的元数据存储。请在 hive-site.xml 中补充:
<property>
<name>javax.jdo.option.ConnectionURL</name>
<value>jdbc:derby:;databaseName=____;create=true</value>
</property>
<property>
<name>hive.metastore.warehouse.dir</name>
<value>____</value>
</property>
答案:____ ____ (每空 1 分,共 2 分)
任务三:搭建实时处理平台
题目 3.1
在 Kafka 配置中,需要设置正确的 Zookeeper 连接。请在 server.properties 中补充:
# Kafka 配置
broker.id=0
listeners=____://localhost:9092
log.dirs=____
zookeeper.connect=____
答案:____ ____ ____ (每空 1 分,共 3 分)
题目 3.2
在 Kafka 主题创建中,需要设置正确的分区和副本数。请补充:
# 创建用户行为主题
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic user-behavior --partitions ____ --replication-factor ____
答案:____ ____ (每空 1 分,共 2 分)
模块二:数据处理(25%)
任务四:数据抽取
题目 4.1
在数据生成器中,需要修复 JSON 序列化问题。请在 data_generator.py 中补充:
# Kafka 生产者配置
self.producer = KafkaProducer(
bootstrap_servers=KAFKA_CONFIG['bootstrap_servers'],
value_serializer=lambda v: json.dumps(v, ensure_ascii=False, default=____).encode('utf-8')
)
答案:____ (1 分)
题目 4.2
在用户行为数据生成中,需要添加正确的行为类型权重。请在 data_generator.py 中补充:
def generate_user_behavior(self):
user = random.choice(self.users)
product = random.choice(self.products)
session = self._get_or_create_session(user.user_id)
# 行为类型权重
behavior_weights = {
BehaviorType.VIEW: ____,
BehaviorType.CLICK: ____,
BehaviorType.ADD_TO_CART: ____,
BehaviorType.REMOVE_FROM_CART: ____,
BehaviorType.SEARCH: ____,
BehaviorType.FAVORITE: 0.05,
BehaviorType.SHARE: ____,
BehaviorType.REVIEW: 0.05
}
答案:____ ____ ____ ____ ____ ____ (每空 1 分,共 6 分)
题目 4.3
在订单数据生成中,需要添加正确的订单状态选择。请补充:
def generate_order(self):
user = random.choice(self.users)
if random.random() > DATA_GENERATION_CONFIG['order_probability']:
return None
# 选择 1-5 个商品
num_items = random.randint(1, 5)
selected_products = random.sample(self.products, num_items)
order_id = str(uuid.uuid4())
order_time = datetime.now()
# 计算订单金额
total_amount = 0
order_items = []
for product in selected_products:
quantity = random.randint(1, 3)
unit_price = product.price
total_price = unit_price * quantity
total_amount += total_price
order_item = OrderItem(
item_id=str(uuid.uuid4()),
order_id=order_id,
product_id=product.product_id,
quantity=quantity,
unit_price=unit_price,
total_price=total_price
)
order_items.append(order_item)
# 计算优惠和最终金额
discount_amount = total_amount * random.uniform(0, 0.3)
shipping_fee = random.uniform(0, 20) if total_amount < ____ else 0
final_amount = total_amount - discount_amount + shipping_fee
order = Order(
user_id=user.user_id,
order_id=order_id,
order_time=order_time,
total_amount=round(total_amount, 2),
discount_amount=round(discount_amount, 2),
final_amount=round(final_amount, 2),
status=random.choice(list(____)),
payment_method=random.choice(self.payment_methods),
shipping_address=f"{random.choice(self.cities)}{random.randint(1, 10)}区{random.randint(1, 100)}号",
shipping_fee=round(shipping_fee, 2),
coupon_code=random.choice([None, f"COUPON_{random.randint(1000, 9999)}"]),
notes=random.choice([None, "请尽快发货", "包装精美", "货到付款"])
)
return order, order_items
答案:____ ____ (每空 1 分,共 2 分)
题目 4.4
在数据生成器初始化中,需要添加正确的设备类型列表。请补充:
def __init__(self):
____
self.producer = KafkaProducer(
bootstrap_servers=KAFKA_CONFIG['bootstrap_servers'],
value_serializer=lambda v: json.dumps(v, ensure_ascii=False, default=str).encode('utf-8')
)
____
# 初始化基础数据
self.cities = ['北京', '上海', '广州', '深圳', '杭州', '南京', '成都', '武汉', '西安', '重庆']
self.device_types = ____
self.browsers = ['Chrome', 'Firefox', 'Safari', 'Edge', 'Opera']
self.payment_methods = ['alipay', 'wechat', 'credit_card', 'debit_card', 'bank_transfer']
self._initialize_data()
答案:____ ____ ____ (每空 1 分,共 3 分)
任务五:数据清洗
题目 5.1
在数据处理器中,需要添加正确的 Kafka 消费者配置。请在 data_processor.py 中补充:
def __init__(self):
self.consumer_behavior = KafkaConsumer(
KAFKA_CONFIG['user_behavior_topic'],
bootstrap_servers=KAFKA_CONFIG['bootstrap_servers'],
auto_offset_reset=____,
enable_auto_commit=____,
group_id=____,
value_deserializer=lambda x: json.loads(x.decode('utf-8'))
)
答案:____ ____ ____ (每空 1 分,共 3 分)
题目 5.2
在数据清洗中,需要添加正确的异常值处理。请补充:
def process_behavior_message(self, message):
try:
behavior = UserBehavior(**message.value)
# 数据清洗:去除异常值
if behavior.duration and (behavior.duration < 0 or behavior.duration > ____):
return # 跳过异常时长数据
if behavior.properties and behavior.properties.get('time_on_page', 0) < 0:
behavior.properties['time_on_page'] = ____
self.behavior_data.append(behavior.model_dump())
# 保存到 CSV
df = pd.DataFrame([behavior.model_dump(mode='json')])
df.to_csv(self.processed_data_path, mode=____, header=False, index=False)
except Exception as e:
print(f"Error processing behavior message: {e}")
答案:____ ____ ____ (每空 1 分,共 3 分)
题目 5.3
在数据分析中,需要添加正确的聚合统计。请补充:
def _perform_analysis(self):
try:
df = pd.read_csv(self.processed_data_path)
if df.empty:
return
# 时间特征提取
df['timestamp'] = pd.____(df['timestamp'])
df['hour'] = df['timestamp'].dt.____
df['day_of_week'] = df['timestamp'].dt.____
# 小时行为统计
hourly_behavior = df.groupby('____')['behavior_type'].____().to_dict()
# 设备类型统计
device_distribution = df.groupby('device_type')['behavior_type'].agg(
lambda x: x.____[0] if not x.mode().empty else 'unknown'
).to_dict()
# 用户活跃度分析
user_activity = df.groupby('user_id')['behavior_type'].____().to_dict()
analysis_results = {
"timestamp": datetime.now().isoformat(),
"hourly_behavior": hourly_behavior,
"device_distribution": device_distribution,
"user_activity": user_activity
}
with open(self.analysis_data_path, 'w') as f:
json.dump(analysis_results, f, ensure_ascii=False, indent=4)
except Exception as e:
print(f"Error performing analysis: {e}")
答案:____ ____ ____ ____ ____ ____ ____ (每空 1 分,共 7 分)
模块三:数据挖掘(15%)
任务六:特征工程
题目 6.1
在机器学习分析器中,需要添加正确的特征工程。请在 ml_analyzer.py 中补充:
def _get_user_data(self):
try:
df_behaviors = pd.read_csv(os.path.join(DATA_DIR, 'user_behaviors.csv'))
df_orders = pd.read_csv(os.path.join(DATA_DIR, 'order_events.csv'))
df_profiles = pd.read_csv(os.path.join(DATA_DIR, 'user_profiles.csv'))
# 转换时间戳
df_behaviors['timestamp'] = pd.to_datetime(df_behaviors['timestamp'])
df_orders['order_date'] = pd.to_datetime(df_orders['order_date'])
df_profiles['last_active'] = pd.____(df_profiles['last_active'])
# 聚合行为数据
user_activity = df_behaviors.groupby('user_id').agg(
total_views=('event_type', lambda x: (x == 'view').sum()),
total_clicks=('event_type', lambda x: (x == 'click').sum()),
total_add_to_cart=('event_type', lambda x: (x == 'add_to_cart').sum()),
total_purchases=('event_type', lambda x: (x == 'purchase').sum()),
last_behavior_time=('timestamp', 'max')
).reset_index()
# 聚合订单数据
user_orders = df_orders.groupby('user_id').agg(
total_order_count=('order_id', 'count'),
total_spent=('total_amount', 'sum'),
last_purchase_date=('order_date', 'max')
).reset_index()
# 合并数据
user_data = pd.merge(df_profiles, user_activity, on='user_id', how='left')
user_data = pd.merge(user_data, user_orders, on='user_id', how='left')
user_data = user_data.fillna(0) # 填充 NaN
# 特征工程
user_data['days_since_last_active'] = (datetime.now() - user_data['last_active']).dt.days
user_data['days_since_last_purchase'] = (datetime.now() - user_data['last_purchase_date']).dt.____
user_data['days_since_last_purchase'] = user_data['days_since_last_purchase'].fillna(user_data['days_since_last_active'])
# 定义流失标签
user_data['churn'] = (user_data['days_since_last_active'] > ____).astype(int)
return user_data
except FileNotFoundError:
print("Required data files not found for ML analysis.")
return pd.DataFrame()
except Exception as e:
print(f"Error loading or processing user data for ML: {e}")
return pd.DataFrame()
答案:____ ____ ____ (每空 1 分,共 3 分)
题目 6.2
在用户分群分析中,需要添加正确的聚类算法。请补充:
def analyze_user_segmentation(self, n_clusters=4):
user_data = self._get_user_data()
if user_data.empty:
return ()
features = user_data[['total_views', 'total_clicks', 'total_purchases', 'total_spent', 'days_since_last_active']]
if self.scaler is None:
self.scaler = StandardScaler()
features_scaled = self.scaler.____(features)
else:
features_scaled = self.scaler.transform(features)
if self.user_segmentation_model is None:
self.user_segmentation_model = KMeans(n_clusters=n_clusters, random_state=____, n_init=10)
self.user_segmentation_model.____(features_scaled)
else:
____
self._save_models()
user_data['segment'] = self.user_segmentation_model.____(features_scaled)
segment_analysis = user_data.groupby('segment').agg(
count=('user_id', 'count'),
avg_views=('total_views', 'mean'),
avg_clicks=('total_clicks', 'mean'),
avg_purchases=('total_purchases', 'mean'),
avg_spent=('total_spent', 'mean'),
avg_days_inactive=('days_since_last_active', 'mean')
).to_dict('index')
return user_data[['user_id', 'segment']].to_dict('records'), segment_analysis
答案:____ ____ ____ ____ (每空 1 分,共 4 分)
任务七:模型训练
题目 7.1
在流失预测模型训练中,需要添加正确的数据分割。请补充:
def train_churn_prediction_model(self):
user_data = self._get_user_data()
if user_data.empty or 'churn' not in user_data.columns or len(user_data['churn'].unique()) < 2:
print("Not enough data or churn labels for training churn prediction model.")
return False
features = user_data[['total_views', 'total_clicks', 'total_purchases', 'total_spent', 'days_since_last_active', 'days_since_last_purchase']]
labels = user_data['churn']
X_train, X_test, y_train, y_test = train_test_split(
features, labels, test_size=____, random_state=____
)
if self.scaler is None:
self.scaler = StandardScaler()
X_train_scaled = self.scaler.____(X_train)
X_test_scaled = self.scaler.transform(X_test)
else:
X_train_scaled = self.scaler.transform(X_train)
X_test_scaled = self.scaler.transform(X_test)
self.churn_prediction_model = RandomForestClassifier(
n_estimators=____,
random_state=42
)
self.churn_prediction_model.____(X_train_scaled, y_train)
y_pred = self.churn_prediction_model.____(X_test_scaled)
accuracy = accuracy_score(y_test, y_pred)
print(f"Churn prediction model trained with accuracy: {accuracy}")
self._save_models()
return True
答案:____ ____ ____ ____ ____ ____ (每空 1 分,共 6 分)
题目 7.2
在商品推荐算法中,需要添加正确的推荐逻辑。请补充:
def recommend_products(self, user_id: str, num_recommendations: int = 5):
user_data = self._get_user_data()
user_profile_row = user_data[user_data['user_id'] == user_id]
if user_profile_row.empty:
return {"user_id": user_id, "recommendations": [], "message": "User not found"}
try:
df_products = pd.read_csv(os.path.join(DATA_DIR, 'product_events.csv'))
except FileNotFoundError:
print("Product data file not found for recommendations.")
return {"user_id": user_id, "recommendations": [], "message": "Product data not available"}
favorite_categories = eval(user_profile_row['favorite_categories'].iloc[0])
recommended_products = []
for category in favorite_categories:
category_products = df_products[df_products['category'] == category].nlargest(
num_recommendations,
'____'
)
recommended_products.extend(category_products['product_id'].tolist())
# 如果推荐不足,用高评分商品补充
if len(recommended_products) < num_recommendations:
top_products = df_products.nlargest(num_recommendations, '____')['product_id'].tolist()
for prod_id in top_products:
if prod_id not in recommended_products:
recommended_products.append(prod_id)
if len(recommended_products) >= num_recommendations:
break
return {"user_id": user_id, "recommendations": recommended_products[:num_recommendations]}
答案:____ ____ (每空 1 分,共 2 分)
模块四:数据采集与实时计算(20%)
任务八:数据采集
题目 8.1
在 Kafka 主题管理中,需要添加正确的主题创建逻辑。请在 start_project.sh 中补充:
# 检查 Kafka 主题是否存在,如果不存在则创建
KAFKA_TOPICS=("user-behavior" "order-events" "product-events" "user-profile")
for TOPIC in "${KAFKA_TOPICS[@]}"; do
EXISTS=$(/home/zkpk/bigdata/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list | grep -w "$TOPIC")
if [ -z "$EXISTS" ]; then
echo "创建 Kafka 主题: $TOPIC"
/home/zkpk/bigdata/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic "$TOPIC" --partitions ____ --replication-factor ____
else
echo "Kafka 主题 '$TOPIC' 已存在。"
fi
done
答案:____ ____ (每空 1 分,共 2 分)
题目 8.2
在数据生成器启动中,需要添加正确的环境变量设置。请补充:
# 启动数据生成器
echo "启动用户行为数据生成器..."
____ "$PROJECT_HOME/src/data_generator.py" > "$LOGS_DIR/data_generator.log" 2>&1 &
echo $! > "$PID_DIR/____"
答案:____ ____ (每空 1 分,共 2 分)
任务九:数据处理
题目 9.1
在 Web 应用中,需要添加正确的 Kafka 消费者线程。请在 app.py 中补充:
def consume_kafka_messages():
global realtime_metrics
consumer_behavior = KafkaConsumer(
KAFKA_CONFIG['____'],
bootstrap_servers=KAFKA_CONFIG['bootstrap_servers'],
auto_offset_reset='____',
enable_auto_commit=____,
group_id='fastapi-behavior-consumer',
value_deserializer=lambda x: json.loads(x.decode('utf-8'))
)
print("开始消费 Kafka 数据...")
for message in consumer_behavior.poll(timeout_ms=100).values():
try:
for msg in message:
behavior = UserBehavior(**msg.value)
realtime_behaviors.append(behavior.model_dump())
# 更新实时指标
realtime_metrics["behavior_distribution"][behavior.event_type] += 1
realtime_metrics["device_distribution"][behavior.device_type] += 1
realtime_metrics["hourly_activity"][behavior.timestamp.hour] += 1
if behavior.event_type == 'view':
realtime_metrics["total_views"] += 1
elif behavior.event_type == 'click':
realtime_metrics["total_clicks"] += 1
elif behavior.event_type == 'add_to_cart':
realtime_metrics["total_add_to_cart"] += 1
elif behavior.event_type == 'purchase':
realtime_metrics["total_purchases"] += 1
except Exception as e:
print(f"Error processing behavior message in FastAPI: {e}, Message: {msg.value}")
答案:____ ____ ____ (每空 1 分,共 3 分)
题目 9.2
在实时数据 API 中,需要添加正确的数据返回格式。请补充:
@app.get("/api/realtime_metrics")
async def get_realtime_metrics():
# 确保小时活动覆盖所有 24 小时
full_hourly_activity = ((h, realtime_metrics["hourly_activity"].get(h, 0)) for h in range(____))
return {
"total_views": realtime_metrics["____"],
"total_clicks": realtime_metrics["total_clicks"],
"total_add_to_cart": realtime_metrics["total_add_to_cart"],
"total_purchases": realtime_metrics["____"],
"total_orders": realtime_metrics["total_orders"],
"total_revenue": round(realtime_metrics["total_revenue"], 2),
"behavior_distribution": dict(realtime_metrics["behavior_distribution"]),
"category_distribution": dict(realtime_metrics["category_distribution"]),
"device_distribution": dict(realtime_metrics["device_distribution"]),
"hourly_activity": dict(sorted(____))
}
答案:____ ____ ____ ____ (4 分)
题目 9.3
在机器学习分析 API 中,需要添加正确的用户分群接口。请补充:
@app.get("/api/user_segmentation")
async def get_user_segmentation():
try:
user_segments, segment_summary = ml_analyzer.____()
return {"user_segments": user_segments, "segment_summary": segment_summary}
except Exception as e:
raise HTTPException(status_code=500, detail=f"User segmentation failed: {e}")
@app.get("/api/churn_prediction/{user_id}")
async def get_churn_prediction(user_id: str):
try:
churn_risk = ml_analyzer.____(user_id)
return churn_risk
except Exception as e:
raise HTTPException(status_code=500, detail=f"Churn prediction failed for user {user_id}: {e}")
@app.get("/api/product_recommendations/{user_id}")
async def get_product_recommendations(user_id: str):
try:
recommendations = ml_analyzer.____(user_id)
return recommendations
except Exception as e:
raise HTTPException(status_code=500, detail=f"Product recommendations failed for user {user_id}: {e}")
答案:____ ____ ____ (每空 1 分,共 3 分)
题目 9.4
在数据流处理中,需要添加正确的实时统计更新。请补充:
def update_traffic_stats(data):
global realtime_metrics
____
behavior_type = data.get('behavior_type', "unknown")
realtime_metrics["behavior_distribution"][behavior_type] += 1
# 更新设备分布统计
device_type = data.get('device_type', "unknown")
realtime_metrics["device_distribution"][device_type] += 1
# 更新小时活跃度统计
timestamp = data.get('timestamp', datetime.now())
if isinstance(timestamp, str):
timestamp = datetime.____(timestamp)
hour = timestamp.____
realtime_metrics["hourly_activity"][hour] += 1
# 更新总计数
if behavior_type == 'view':
realtime_metrics["total_views"] += 1
elif behavior_type == 'click':
realtime_metrics["total_clicks"] += 1
elif behavior_type == 'add_to_cart':
realtime_metrics["total_add_to_cart"] += 1
elif behavior_type == 'purchase':
realtime_metrics["total_purchases"] += 1
realtime_metrics["total_orders"] += 1
realtime_metrics["total_revenue"] += data.get('total_amount', 0)
答案:____ ____ ____ (3 分)
题目 9.5
在数据持久化中,需要添加正确的数据保存逻辑。请补充:
def save_data_periodically(self):
if self.behavior_data:
df_behavior = pd.DataFrame(self.behavior_data)
df_behavior['timestamp'] = pd.to_datetime(df_behavior['timestamp'])
df_behavior.to_csv(
os.path.join(DATA_DIR, 'user_behaviors.csv'),
mode=____,
header=not os.path.exists(os.path.join(DATA_DIR, 'user_behaviors.csv')),
index=False
)
self.behavior_data = [] # 清空数据
if self.order_data:
df_orders = pd.DataFrame(self.order_data)
df_orders['order_date'] = pd.to_datetime(df_orders['order_date'])
df_orders.to_csv(
os.path.join(DATA_DIR, 'order_events.csv'),
mode=____,
header=not os.path.exists(os.path.join(DATA_DIR, 'order_events.csv')),
index=False
)
self.order_data = [] # 清空数据
if self.user_profile_data:
df_user_profiles = pd.DataFrame(list(self.user_profile_data.values()))
df_user_profiles['last_active'] = pd.to_datetime(df_user_profiles['last_active'])
df_user_profiles.to_csv(
os.path.join(DATA_DIR, 'user_profiles.csv'),
mode=____,
index=False
) # 覆盖保存最新状态
答案:____ ____ ____ (每空 1 分,共 3 分)
模块五:数据可视化(10%)
任务十:制作图表
题目 10.1
在前端 JavaScript 中,需要添加正确的图表初始化。请在 index.html 中补充:
// 初始化图表
function initCharts() {
charts.behaviorDistribution = echarts.init(document.getElementById('behaviorDistributionChart'));
charts.deviceDistribution = echarts.____(document.getElementById('deviceDistributionChart'));
charts.hourlyActivity = echarts.init(document.getElementById('hourlyActivityChart'));
charts.userSegmentation = echarts.init(document.getElementById('userSegmentationChart'));
// 响应式调整
window.addEventListener('resize', function() {
Object.values(charts).forEach(chart => chart.____());
});
}
答案:____ ____ (每空 1 分,共 2 分)
题目 10.2
在用户 behavior 分布图表中,需要添加正确的饼图配置。请补充:
// 更新用户行为分布图表
function updateBehaviorChart(data) {
const behaviorData = Object.entries(data.behavior_distribution).map(([name, value]) => ({
name: name === 'view' ? '预览' :
name === 'click' ? '点击' :
name === 'add_to_cart' ? '加购物车' :
name === 'purchase' ? '购买' :
name === 'favorite' ? '收藏' : '分享',
value: ____
}));
charts.behaviorDistribution.setOption({
title: { text: '用户行为分布', left: 'center', show: false },
tooltip: { trigger: 'item' },
legend: { orient: 'vertical', left: 'left' },
series: [{
name: '行为类型',
type: '____',
radius: '____',
data: behaviorData,
emphasis: {
itemStyle: {
shadowBlur: 10,
shadowOffsetX: 0,
shadowColor: 'rgba(0, 0, 0, 0.5)'
}
}
}]
});
}
答案:____ ____ ____ (每空 1 分,共 3 分)
题目 10.3
在 24 小时活跃度趋势图表中,需要添加正确的时间序列配置。请补充:
// 更新 24 小时活跃度趋势图表
function updateHourlyActivityChart(data) {
const hours = Object.keys(data.hourly_activity).map(hour => `${hour}:00`);
const activityData = Object.values(data.hourly_activity);
charts.hourlyActivity.setOption({
title: { text: '24小时活跃度趋势', left: 'center', show: false },
tooltip: { trigger: 'axis' },
xAxis: { type: 'category', data: hours },
yAxis: { type: 'value', name: '活跃度' },
series: [{
name: '活跃度',
type: '____',
data: activityData,
smooth: ____,
itemStyle: { color: '#91cc75' }
}]
});
}
答案:____ ____ (每空 1 分,共 2 分)
任务十一:分析数据
题目 11.1
在实时监控功能中,需要添加正确的定时更新逻辑。请补充:
// 开始实时监控
function startRealTime() {
if (realTimeInterval) clearInterval(realTimeInterval);
realTimeInterval = setInterval(async () => {
try {
const response = await fetch('____');
const data = await response.json();
if (data.____ === 'success') {
updateBehaviorChart(data);
updateDeviceChart(data);
updateHourlyActivityChart(data);
// 更新指标显示
document.getElementById('totalViews').innerText = data.total_views;
document.getElementById('totalClicks').innerText = data.total_clicks;
document.getElementById('totalPurchases').innerText = data.total_purchases;
document.getElementById('totalRevenue').innerText = data.total_revenue.toFixed(2);
document.getElementById('status').textContent = '状态:实时监控中 - ' + new Date().toLocaleTimeString();
}
} catch (error) {
console.error('实时监控出错:', error);
}
}, ____); // 每 3 秒更新一次
document.getElementById('startBtn').disabled = true;
document.getElementById('stopBtn').disabled = false;
document.getElementById('status').textContent = '状态:实时监控已启动';
}
答案:____ ____ ____ (每空 1 分,共 3 分)
模块六:综合分析任务(5%)
任务十二:分析数据
题目 12.1
在系统健康检查中,需要添加正确的服务状态监控。请在 app.py 中补充:
@app.get("/api/health")
async def health_check():
kafka_status = "healthy"
try:
# 尝试列出主题来检查 Kafka 连接
____
except Exception:
kafka_status = "unhealthy"
ml_model_status = "healthy" if ml_analyzer.user_segmentation_model and ml_analyzer.churn_prediction_model else "uninitialized"
return {
"status": "ok",
"kafka_status": kafka_status,
"ml_models": ml_model_status,
"timestamp": datetime.now()
}
答案:____ (1 分)
题目 12.2
在数据分析 API 中,需要添加正确的统计计算。请补充:
@app.get("/api/behavior/analysis")
async def get_behavior_analysis():
try:
return {
"status": "____",
"data": {
"behavior_distribution": realtime_metrics["behavior_distribution"],
"device_distribution": realtime_metrics["behavior_distribution"],
"hourly_activity": [
{
"hour": f"{i}:00",
"activity": realtime_metrics["hourly_activity"].get(i, 0)
} for i in range(____)
]
}
}
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
答案:____ ____ (每空 1 分,共 2 分)
任务十三:提出方案
题目 13.1
在项目启动脚本中,需要添加正确的服务启动顺序。请在 start_project.sh 中补充:
#!/bin/bash
PROJECT_HOME="/home/zkpk/projects/ecommerce-user-behavior"
PYTHON_ENV="/home/zkpk/Anaconda3/bin/python"
KAFKA_HOME="/home/zkpk/bigdata/kafka"
echo "启动电商用户行为分析平台..."
# 1. 检查大数据环境状态
echo "检查大数据环境状态..."
/home/zkpk/bigdata/check-bigdata.sh
if [ $? -ne 0 ]; then
echo "大数据环境检查失败,请确保所有组件已正确启动。"
exit ____
fi
# 2. 检查并创建 Kafka 主题
echo "检查 Kafka 主题..."
export PATH=$PATH:$KAFKA_HOME/bin
____
# 3. 启动数据生成器
echo "启动用户行为数据生成器..."
nohup $PYTHON_ENV $PROJECT_HOME/src/data_generator.py > $PROJECT_HOME/logs/data_generator.log 2>&1 &
echo $! > $PROJECT_HOME/pids/data_generator.pid
sleep 2
# 4. 启动数据处理器
echo "启动数据处理器..."
nohup $PYTHON_ENV $PROJECT_HOME/src/data_processor.py > $PROJECT_HOME/logs/data_processor.log 2>&1 &
echo $! > $PROJECT_HOME/pids/data_processor.pid
sleep 2
# 5. 启动 Web 服务
echo "启动 Web 服务..."
export PYTHONPATH=$PROJECT_HOME:$PYTHONPATH
nohup $PYTHON_ENV -m uvicorn web.app:app --host 0.0.0.0 --port 8080 > $PROJECT_HOME/logs/web_service.log 2>&1 &
echo $! > $PROJECT_HOME/pids/web_service.pid
sleep 5
echo "平台启动完成!"
echo "Web 界面: http://localhost:8080"
echo "使用 $PROJECT_HOME/stop_project.sh 停止所有服务"
答案:____ ____ (每空 1 分,共 2 分)