利用Flink在大数据领域构建实时推荐系统
利用Flink在大数据领域构建实时推荐系统
关键词:Flink、实时推荐系统、大数据处理、机器学习、流式计算、个性化推荐、分布式系统
摘要:本文深入探讨如何利用Apache Flink构建高性能的实时推荐系统。我们将从推荐系统的基本原理出发,详细分析Flink在实时数据处理中的优势,介绍推荐算法的实现方式,并通过完整的项目案例展示如何构建端到端的实时推荐系统。文章将覆盖数据处理流水线设计、特征工程、模型训练与在线预测等关键环节,同时讨论系统性能优化和扩展性考虑。
1. 背景介绍
1.1 目的和范围
本文旨在为大数据工程师和架构师提供构建实时推荐系统的完整指南。我们将重点讨论:
- Flink在实时推荐系统中的核心作用
- 推荐系统的基本架构和关键组件
- 实时特征计算和模型更新的技术实现
- 系统性能优化和扩展性考虑
1.2 预期读者
本文适合以下读者:
- 大数据工程师:希望了解如何利用Flink构建实时数据处理系统
- 机器学习工程师:需要将模型部署到实时生产环境
- 系统架构师:设计高可用、高性能的推荐系统架构
- 技术决策者:评估实时推荐系统的技术选型和实施方案
1.3 文档结构概述
本文采用从理论到实践的递进结构:
- 首先介绍推荐系统和Flink的基本概念
- 然后深入分析推荐系统的核心算法和数学模型
- 接着通过完整项目案例展示实现细节
- 最后讨论实际应用中的挑战和解决方案
1.4 术语表
1.4.1 核心术语定义
- Flink:Apache Flink是一个分布式流处理框架,支持有状态的计算和精确一次处理语义
- 实时推荐系统:能够在毫秒到秒级延迟内响应用户行为并提供个性化推荐的系统
- 协同过滤:基于用户历史行为数据发现用户偏好的推荐算法
- 特征工程:将原始数据转换为机器学习模型可理解的特征的过程
- 模型服务:将训练好的模型部署为可实时响应预测请求的服务
1.4.2 相关概念解释
- 流批一体:Flink的核心特性,同一套代码可以处理流数据和批数据
- CEP(复杂事件处理):从事件流中识别特定模式的技术
- Embedding:将高维稀疏特征映射到低维连续向量空间的技术
- A/B测试:比较两种不同算法或系统效果的实验方法
1.4.3 缩略词列表
- CF: Collaborative Filtering (协同过滤)
- MF: Matrix Factorization (矩阵分解)
- ALS: Alternating Least Squares (交替最小二乘法)
- CTR: Click-Through Rate (点击率)
- ROC: Receiver Operating Characteristic (受试者工作特征曲线)
- AUC: Area Under Curve (曲线下面积)
2. 核心概念与联系
2.1 实时推荐系统架构
2.2 Flink在推荐系统中的角色
Flink在实时推荐系统中扮演着核心角色:
- 数据接入层:从Kafka、数据库等源实时消费数据
- 流式处理引擎:执行窗口聚合、特征计算等操作
- 状态管理:维护用户最近行为、会话状态等
- 模型服务桥梁:连接特征工程和模型预测服务
2.3 推荐系统关键组件交互
3. 核心算法原理 & 具体操作步骤
3.1 实时推荐算法概述
实时推荐系统通常结合以下算法:
- 基于内容的推荐:利用物品特征匹配用户兴趣
- 协同过滤:基于用户-物品交互矩阵发现相似性
- 深度学习模型:如Wide & Deep、DeepFM等
- 上下文感知推荐:考虑时间、位置等上下文因素
3.2 Flink实现协同过滤算法
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, DataTypes
from pyflink.table.descriptors import Schema, Kafka, Json
# 初始化Flink环境
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# 定义Kafka数据源
t_env.connect(Kafka()
.version("universal")
.topic("user_behavior")
.start_from_earliest()
.property("zookeeper.connect", "localhost:2181")
.property("bootstrap.servers", "localhost:9092")) \
.with_format(Json()
.fail_on_missing_field(True)
.schema(DataTypes.ROW([
DataTypes.FIELD("user_id", DataTypes.BIGINT()),
DataTypes.FIELD("item_id", DataTypes.BIGINT()),
DataTypes.FIELD("behavior", DataTypes.STRING()),
DataTypes.FIELD("timestamp", DataTypes.BIGINT())
]))) \
.with_schema(Schema()
.field("user_id", DataTypes.BIGINT())
.field("item_id", DataTypes.BIGINT())
.field("behavior", DataTypes.STRING())
.field("timestamp", DataTypes.BIGINT())) \
.create_temporary_table("user_behavior")
# 计算用户-物品交互矩阵
t_env.sql_query("""
SELECT
user_id,
item_id,
COUNT(*) as interaction_count,
MAX(timestamp) as last_interaction_time
FROM user_behavior
WHERE behavior = 'click' OR behavior = 'purchase'
GROUP BY user_id, item_id
""").to_retract_stream().print()
# 执行作业
env.execute("Real-time Collaborative Filtering")
3.3 实时矩阵分解算法
矩阵分解是推荐系统的经典算法,实时版本需要考虑增量更新:
- **交替最小二乘法(ALS)**的在线变体
- 使用Flink状态维护用户和物品的隐向量
- 通过窗口聚合实现增量学习
import numpy as np
from pyflink.datastream import MapFunction, RuntimeContext
from pyflink.common.typeinfo import Types
from pyflink.common import Row
class OnlineALS(MapFunction):
def __init__(self, rank, lambda_):
self.rank = rank # 隐向量维度
self.lambda_ = lambda_ # 正则化系数
def open(self, runtime_context: RuntimeContext):
# 从状态中加载用户和物品矩阵
self.user_state = runtime_context.get_state(
"user_matrix",
Types.PICKLED_BYTE_ARRAY())
self.item_state = runtime_context.get_state(
"item_matrix",
Types.PICKLED_BYTE_ARRAY())
# 初始化随机矩阵
if not self.user_state.value():
self.user_matrix = {}
else:
self.user_matrix = self.user_state.value()
if not self.item_state.value():
self.item_matrix = {}
else:
self.item_matrix = self.item_state.value()
def map(self, value):
user_id = value[0]
item_id = value[1]
rating = value[2]
# 如果用户/物品不存在,初始化隐向量
if user_id not in self.user_matrix:
self.user_matrix[user_id] = np.random.randn(self.rank)
if item_id not in self.item_matrix:
self.item_matrix[item_id] = np.random.randn(self.rank)
# 获取当前向量
user_vec = self.user_matrix[user_id]
item_vec = self.item_matrix[item_id]
# 计算预测误差
prediction = np.dot(user_vec, item_vec)
error = rating - prediction
# 更新用户向量
item_vec_norm = np.dot(item_vec, item_vec)
user_vec_new = user_vec + 0.01 * (
error * item_vec - self.lambda_ * user_vec)
# 更新物品向量
user_vec_norm = np.dot(user_vec, user_vec)
item_vec_new = item_vec + 0.01 * (
error * user_vec - self.lambda_ * item_vec)
# 保存更新
self.user_matrix[user_id] = user_vec_new
self.item_matrix[item_id] = item_vec_new
# 更新状态
self.user_state.update(self.user_matrix)
self.item_state.update(self.item_matrix)
return Row.of(user_id, item_id, float(prediction))
# 使用示例
als = OnlineALS(rank=10, lambda_=0.1)
data_stream.map(als).returns(Types.ROW([Types.LONG(), Types.LONG(), Types.FLOAT()]))
4. 数学模型和公式 & 详细讲解 & 举例说明
4.1 推荐系统基础数学模型
推荐系统的核心是预测用户对物品的偏好程度,通常表示为:
r^ui=f(u,i∣Θ) \hat{r}_{ui} = f(u,i|\Theta) r^ui=f(u,i∣Θ)
其中:
- r^ui\hat{r}_{ui}r^ui 是用户u对物品i的预测评分
- fff 是预测函数
- Θ\ThetaΘ 是模型参数
4.2 矩阵分解模型
矩阵分解将用户-物品交互矩阵R分解为用户矩阵U和物品矩阵V:
R≈U×VT R \approx U \times V^T R≈U×VT
优化目标是最小化以下损失函数:
minU,V∑(u,i)∈K(rui−uuTvi)2+λ(∥uu∥2+∥vi∥2) \min_{U,V} \sum_{(u,i)\in K} (r_{ui} - u_u^T v_i)^2 + \lambda (\|u_u\|^2 + \|v_i\|^2) U,Vmin(u,i)∈K∑(rui−uuTvi)2+λ(∥uu∥2+∥vi∥2)
其中:
- KKK 是已知评分的集合
- λ\lambdaλ 是正则化系数
- uuu_uuu 是用户u的隐向量
- viv_ivi 是物品i的隐向量
4.3 实时增量学习
在实时场景下,我们需要增量更新模型参数。对于每个新到达的评分(u,i,rui)(u,i,r_{ui})(u,i,rui),参数更新规则为:
用户向量更新:
uu←uu+γ(eui⋅vi−λuu) u_u \leftarrow u_u + \gamma (e_{ui} \cdot v_i - \lambda u_u) uu←uu+γ(eui⋅vi−λuu)
物品向量更新:
vi←vi+γ(eui⋅uu−λvi) v_i \leftarrow v_i + \gamma (e_{ui} \cdot u_u - \lambda v_i) vi←vi+γ(eui⋅uu−λvi)
其中:
- eui=rui−uuTvie_{ui} = r_{ui} - u_u^T v_ieui=rui−uuTvi 是预测误差
- γ\gammaγ 是学习率
4.4 多目标优化
现代推荐系统通常需要优化多个目标,如点击率(CTR)和停留时长。我们可以使用多任务学习框架:
L=αLCTR+(1−α)LWatchTime+λ∥Θ∥2 L = \alpha L_{CTR} + (1-\alpha)L_{WatchTime} + \lambda \|\Theta\|^2 L=αLCTR+(1−α)LWatchTime+λ∥Θ∥2
其中α\alphaα是平衡两个目标的超参数。
5. 项目实战:代码实际案例和详细解释说明
5.1 开发环境搭建
5.1.1 环境要求
- Java 8+
- Python 3.7+
- Apache Flink 1.13+
- Kafka 2.8+
- Redis 6.0+ (用于特征缓存)
5.1.2 依赖安装
# 安装PyFlink
pip install apache-flink
# 安装机器学习相关库
pip install numpy scikit-learn tensorflow
# 安装特征存储客户端
pip install redis
5.2 源代码详细实现和代码解读
5.2.1 实时特征工程实现
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, DataTypes
from pyflink.table.descriptors import Schema, Kafka, Json
from pyflink.table.window import Tumble
import json
import redis
class RealTimeFeatureEngineering:
def __init__(self):
self.env = StreamExecutionEnvironment.get_execution_environment()
self.table_env = StreamTableEnvironment.create(self.env)
self.redis_client = redis.Redis(host='localhost', port=6379, db=0)
def setup_kafka_source(self):
# 配置Kafka数据源
self.table_env.connect(Kafka()
.version("universal")
.topic("user_events")
.start_from_earliest()
.property("bootstrap.servers", "localhost:9092")) \
.with_format(Json()
.fail_on_missing_field(True)
.schema(DataTypes.ROW([
DataTypes.FIELD("user_id", DataTypes.STRING()),
DataTypes.FIELD("item_id", DataTypes.STRING()),
DataTypes.FIELD("event_type", DataTypes.STRING()),
DataTypes.FIELD("timestamp", DataTypes.BIGINT()),
DataTypes.FIELD("extra", DataTypes.STRING())
]))) \
.with_schema(Schema()
.field("user_id", DataTypes.STRING())
.field("item_id", DataTypes.STRING())
.field("event_type", DataTypes.STRING())
.field("timestamp", DataTypes.BIGINT())
.field("extra", DataTypes.STRING())) \
.create_temporary_table("user_events")
def compute_user_features(self):
# 计算用户近期行为特征
user_features = self.table_env.sql_query("""
SELECT
user_id,
COUNT(CASE WHEN event_type = 'click' THEN 1 END) AS click_count_1h,
COUNT(CASE WHEN event_type = 'purchase' THEN 1 END) AS purchase_count_1h,
COUNT(CASE WHEN event_type = 'like' THEN 1 END) AS like_count_1h,
HOP_ROWTIME(rowtime, INTERVAL '5' SECOND, INTERVAL '1' HOUR) AS feature_time
FROM user_events
GROUP BY
HOP(rowtime, INTERVAL '5' SECOND, INTERVAL '1' HOUR),
user_id
""")
# 将特征写入Redis
def sink_to_redis(row):
user_id = row["user_id"]
features = {
"click_count_1h": row["click_count_1h"],
"purchase_count_1h": row["purchase_count_1h"],
"like_count_1h": row["like_count_1h"],
"last_update": row["feature_time"].strftime("%Y-%m-%d %H:%M:%S")
}
self.redis_client.hset(f"user:{user_id}:features", mapping=features)
user_features.to_append_stream().map(sink_to_redis)
def compute_item_features(self):
# 计算物品实时热度特征
item_features = self.table_env.sql_query("""
SELECT
item_id,
COUNT(*) AS interaction_count_1h,
COUNT(DISTINCT user_id) AS unique_users_1h,
TUMBLE_ROWTIME(rowtime, INTERVAL '1' HOUR) AS feature_time
FROM user_events
GROUP BY
TUMBLE(rowtime, INTERVAL '1' HOUR),
item_id
""")
def sink_item_features(row):
item_id = row["item_id"]
features = {
"interaction_count_1h": row["interaction_count_1h"],
"unique_users_1h": row["unique_users_1h"],
"last_update": row["feature_time"].strftime("%Y-%m-%d %H:%M:%S")
}
self.redis_client.hset(f"item:{item_id}:features", mapping=features)
item_features.to_append_stream().map(sink_item_features)
def run(self):
self.setup_kafka_source()
self.compute_user_features()
self.compute_item_features()
self.env.execute("Real-time Feature Engineering")
if __name__ == "__main__":
feature_engineering = RealTimeFeatureEngineering()
feature_engineering.run()
5.2.2 实时推荐服务实现
from pyflink.datastream import StreamExecutionEnvironment, FlatMapFunction
from pyflink.common.typeinfo import Types
from pyflink.common import Row
import numpy as np
import redis
import json
from typing import List, Tuple
class RealTimeRecommender(FlatMapFunction):
def __init__(self):
self.redis_client = redis.Redis(host='localhost', port=6379, db=0)
self.user_model = {} # 用户模型缓存
self.item_model = {} # 物品模型缓存
def load_user_model(self, user_id: str) -> np.ndarray:
# 从Redis加载用户模型
if user_id not in self.user_model:
model_data = self.redis_client.get(f"user:{user_id}:model")
if model_data:
self.user_model[user_id] = np.frombuffer(model_data, dtype=np.float32)
else:
# 初始化随机向量
self.user_model[user_id] = np.random.rand(10).astype(np.float32)
return self.user_model[user_id]
def load_item_model(self, item_id: str) -> np.ndarray:
# 从Redis加载物品模型
if item_id not in self.item_model:
model_data = self.redis_client.get(f"item:{item_id}:model")
if model_data:
self.item_model[item_id] = np.frombuffer(model_data, dtype=np.float32)
else:
# 初始化随机向量
self.item_model[item_id] = np.random.rand(10).astype(np.float32)
return self.item_model[item_id]
def get_user_features(self, user_id: str) -> dict:
# 获取用户实时特征
features = self.redis_client.hgetall(f"user:{user_id}:features")
return {
"click_count_1h": int(features.get(b"click_count_1h", 0)),
"purchase_count_1h": int(features.get(b"purchase_count_1h", 0)),
"like_count_1h": int(features.get(b"like_count_1h", 0))
}
def get_item_features(self, item_id: str) -> dict:
# 获取物品实时特征
features = self.redis_client.hgetall(f"item:{item_id}:features")
return {
"interaction_count_1h": int(features.get(b"interaction_count_1h", 0)),
"unique_users_1h": int(features.get(b"unique_users_1h", 0))
}
def predict_score(self, user_id: str, item_id: str) -> float:
# 计算预测分数
user_vec = self.load_user_model(user_id)
item_vec = self.load_item_model(item_id)
# 获取实时特征
user_features = self.get_user_features(user_id)
item_features = self.get_item_features(item_id)
# 简单加权分数
base_score = np.dot(user_vec, item_vec)
time_decay = 0.1 * user_features["click_count_1h"] + 0.3 * user_features["purchase_count_1h"]
popularity = 0.01 * item_features["interaction_count_1h"]
return float(base_score + time_decay + popularity)
def recommend_items(self, user_id: str, candidate_items: List[str], top_n: int = 10) -> List[Tuple[str, float]]:
# 为指定用户推荐物品
scored_items = []
for item_id in candidate_items:
score = self.predict_score(user_id, item_id)
scored_items.append((item_id, score))
# 按分数排序并返回topN
scored_items.sort(key=lambda x: x[1], reverse=True)
return scored_items[:top_n]
def flat_map(self, value):
# 处理推荐请求
request = json.loads(value)
user_id = request["user_id"]
candidate_items = request["candidate_items"]
top_n = request.get("top_n", 10)
# 生成推荐结果
recommendations = self.recommend_items(user_id, candidate_items, top_n)
# 返回推荐结果
for item_id, score in recommendations:
yield Row.of(user_id, item_id, float(score))
# 使用示例
env = StreamExecutionEnvironment.get_execution_environment()
recommender = RealTimeRecommender()
# 假设从Kafka读取推荐请求
requests = env.add_source(KafkaSource...) # 省略Kafka源配置
recommendations = requests.flat_map(recommender) \
.returns(Types.ROW([Types.STRING(), Types.STRING(), Types.FLOAT()]))
# 将推荐结果写入Kafka或其他存储
recommendations.add_sink(...)
env.execute("Real-time Recommendation Service")
5.3 代码解读与分析
5.3.1 特征工程模块
- 数据源配置:通过Flink SQL API连接Kafka数据源
- 窗口计算:使用TUMBLE和HOP窗口函数计算时间窗口内的统计特征
- 特征存储:将计算好的特征实时写入Redis供推荐服务使用
- 状态管理:利用Flink的状态机制保证特征计算的准确性和一致性
5.3.2 推荐服务模块
- 模型加载:从Redis加载预训练的用户和物品Embedding
- 特征融合:结合静态模型和实时特征进行综合评分
- 候选集筛选:从海量物品中高效筛选TopN推荐结果
- 流式处理:以事件驱动的方式实时响应推荐请求
5.3.3 系统优势
- 低延迟:毫秒级的推荐响应时间
- 实时性:能够捕捉用户最新行为并立即反映在推荐结果中
- 可扩展性:Flink的分布式架构支持水平扩展
- 一致性:精确一次处理语义保证数据处理的准确性
6. 实际应用场景
6.1 电商个性化推荐
- 实时个性化首页:根据用户实时浏览行为调整商品展示顺序
- 购物车推荐:基于当前购物车内容推荐相关商品
- 限时抢购推荐:结合库存和用户偏好实时调整推荐策略
6.2 内容平台推荐
- 新闻资讯推荐:根据用户点击流实时调整新闻排序
- 视频推荐:基于观看时长和互动行为优化推荐内容
- 社交feed流:动态调整内容展示顺序和多样性
6.3 游戏道具推荐
- 游戏内商城:根据玩家实时游戏行为和装备情况推荐道具
- 赛季通行证:动态调整奖励物品推荐策略
- 社交推荐:实时推荐可能感兴趣的队友或对手
6.4 金融产品推荐
- 个性化理财产品:根据用户风险偏好和实时市场数据推荐产品
- 信用卡推荐:结合用户消费习惯实时推荐优惠活动
- 保险产品:基于用户画像和生活事件触发推荐
7. 工具和资源推荐
7.1 学习资源推荐
7.1.1 书籍推荐
- 《Stream Processing with Apache Flink》- Fabian Hueske
- 《Recommender Systems: The Textbook》- Charu C. Aggarwal
- 《Real-Time Analytics: Techniques to Analyze and Visualize Streaming Data》- Byron Ellis
7.1.2 在线课程
- Coursera: “Big Data Analysis with Scala and Spark” (含Flink内容)
- Udemy: “Apache Flink for Real-Time Stream Processing”
- edX: “Scalable Machine Learning” (含实时推荐系统内容)
7.1.3 技术博客和网站
- Flink官方博客: https://flink.apache.org/blog/
- Netflix Tech Blog: 实时推荐系统实践
- Airbnb Engineering: 实时个性化技术分享
7.2 开发工具框架推荐
7.2.1 IDE和编辑器
- IntelliJ IDEA (Flink开发最佳选择)
- VS Code with PyFlink插件
- Jupyter Notebook for 原型开发
7.2.2 调试和性能分析工具
- Flink Web UI (内置监控界面)
- Prometheus + Grafana (指标监控)
- JProfiler (性能分析)
7.2.3 相关框架和库
- Apache Kafka (事件流处理)
- Redis (特征缓存)
- TensorFlow/PyTorch (深度学习模型)
- Feast (特征存储)
7.3 相关论文著作推荐
7.3.1 经典论文
- “Matrix Factorization Techniques for Recommender Systems” (Koren et al.)
- “Deep Neural Networks for YouTube Recommendations” (Covington et al.)
- “Real-time Personalization using Embeddings for Search Ranking at Airbnb” (Grbovic et al.)
7.3.2 最新研究成果
- “Streaming Recommender Systems” (ACM SIGIR 2022)
- “Real-Time Recommendation System with Reinforcement Learning” (KDD 2023)
- “Edge-Enabled Real-Time Recommendations” (IEEE IoT Journal 2023)
7.3.3 应用案例分析
- LinkedIn实时推荐系统架构
- Alibaba双11实时个性化实践
- Spotify音乐推荐系统演进
8. 总结:未来发展趋势与挑战
8.1 发展趋势
- 更智能的实时学习:结合强化学习和在线学习的混合方法
- 多模态推荐:整合文本、图像、视频等多模态数据
- 边缘计算:将部分推荐逻辑下放到边缘设备减少延迟
- 可解释推荐:提供实时可解释的推荐理由
- 隐私保护:在实时场景下实现联邦学习和差分隐私
8.2 技术挑战
- 冷启动问题:实时系统中新用户/物品的快速适应
- 数据倾斜:热门内容导致的负载不均衡
- 状态管理:大规模分布式状态的一致性和性能平衡
- 模型漂移:实时数据分布变化导致的模型性能下降
- 系统复杂性:多组件协同的运维复杂度
8.3 实践建议
- 渐进式演进:从批处理逐步过渡到实时处理
- 监控先行:建立完善的指标监控和告警系统
- A/B测试文化:任何变更都通过实验验证效果
- 模块化设计:保持系统各组件的独立性和可替换性
- 性能与效果平衡:在延迟和推荐质量间找到最佳平衡点
9. 附录:常见问题与解答
Q1: Flink和Spark Streaming在实时推荐系统中的选择依据是什么?
A1: 主要考虑因素包括:
- 延迟要求:Flink的毫秒级延迟优于Spark Streaming的秒级延迟
- 状态管理:Flink的状态管理更完善,适合复杂推荐逻辑
- 精确一次语义:两者都支持,但Flink实现更轻量级
- 机器学习集成:Spark MLlib更成熟,但Flink ML在快速发展
Q2: 如何处理实时推荐系统中的数据倾斜问题?
A2: 常用策略有:
- 预分区:对热点键进行特殊分区处理
- 本地聚合:先在本地进行部分聚合再全局聚合
- 倾斜键分离:将热点数据单独处理
- 动态负载均衡:基于实时监控调整资源分配
Q3: 实时推荐系统如何评估效果?
A3: 需要结合多种评估方式:
- 在线指标:实时CTR、转化率、停留时长等
- 离线指标:AUC、NDCG等传统推荐指标
- A/B测试:新旧算法在真实流量上的对比
- 人工评估:定期抽样进行人工质量评估
Q4: 如何保证实时推荐系统的高可用性?
A4: 关键措施包括:
- 多副本:Flink作业和存储系统的多副本部署
- 容错机制:利用Flink的检查点和保存点
- 降级策略:在异常情况下回退到简化版本
- 流量控制:实现熔断和限流机制
Q5: 实时推荐系统如何解决用户隐私问题?
A5: 可采用的隐私保护技术:
- 数据脱敏:去除直接个人标识信息
- 差分隐私:在特征计算中注入可控噪声
- 联邦学习:数据保留在本地只共享模型更新
- 用户授权:明确的数据使用授权和透明化
10. 扩展阅读 & 参考资料
- Apache Flink官方文档: https://flink.apache.org/
- Netflix实时推荐系统架构: https://netflixtechblog.com/
- 《Building Machine Learning Powered Applications》- Emmanuel Ameisen
- ACM RecSys会议论文集
- IEEE Transactions on Knowledge and Data Engineering中相关论文
- Flink社区最佳实践案例
- O’Reilly《Real-Time Analytics》系列文章
- Google Research关于实时机器学习的博客
- LinkedIn Engineering的技术博客
- Alibaba双11技术复盘中的推荐系统章节
更多推荐



所有评论(0)