1. PyFlink ML架构概览
1.1 流式机器学习范式
传统批处理ML与流式ML的核心差异。
# 传统批处理ML vs 流式ML对比
from pyflink.ml import MLPipeline
from pyflink.ml.feature import StandardScaler, VectorAssembler
from pyflink.ml.classification import LogisticRegression
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
# 批处理ML(静态数据集)
def batch_ml_pipeline():
# 加载完整数据集
# 训练模型
# 批量预测
pass
# 流式ML(持续数据流)
def streaming_ml_pipeline(env: StreamExecutionEnvironment):
# 持续数据流
# 增量学习/在线学习
# 实时预测与模型更新
pass
# 架构对比
"""
批处理ML特征:
├── 静态数据集训练
├── 固定模型版本
├── 周期性重训练
├── 高延迟预测
└── 资源密集型
流式ML特征:
├── 持续数据流训练
├── 动态模型演化
├── 实时模型更新
├── 低延迟预测
└── 资源高效
"""
1.2 PyFlink ML组件体系
完整的ML生态系统架构。
from pyflink.ml import Model, Estimator
from pyflink.ml.core import AlgoOperator, Model
from pyflink.ml.api import Transformer, Predictor
# PyFlink ML核心组件
class MLComponents:
"""ML组件分类"""
# 数据预处理
PREPROCESSING = [
'StandardScaler', # 标准化
'MinMaxScaler', # 归一化
'StringIndexer', # 字符串索引
'OneHotEncoder', # 独热编码
'VectorAssembler', # 特征向量组装
'PCA', # 主成分分析
]
# 特征工程
FEATURE_ENGINEERING = [
'PolynomialExpansion', # 多项式扩展
'Bucketizer', # 分桶
'ElementwiseProduct', # 元素乘积
'NGram', # N元语法
]
# 机器学习算法
ALGORITHMS = [
# 分类
'LogisticRegression',
'DecisionTreeClassifier',
'RandomForestClassifier',
'LinearSVC',
# 回归
'LinearRegression',
'DecisionTreeRegressor',
'GBTRegressor',
# 聚类
'KMeans',
'GaussianMixture',
# 推荐
'ALS',
]
# 模型评估与调优
EVALUATION = [
'BinaryClassificationEvaluator',
'MulticlassClassificationEvaluator',
'RegressionEvaluator',
'CrossValidator',
]
2. 实时特征工程
2.1 流式特征提取
实时数据流的特征计算管道。
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, DataTypes
from pyflink.ml.feature import StandardScaler, VectorAssembler, StringIndexer
from pyflink.ml import Pipeline
import pandas as pd
# 创建流处理环境
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# 实时用户行为特征工程
def create_user_behavior_features():
"""实时用户行为特征计算"""
# 1. 创建源数据流(Kafka实时数据)
source_ddl = """
CREATE TABLE user_behavior_stream (
user_id BIGINT,
item_id BIGINT,
behavior_type STRING, — click, purchase, view
behavior_time TIMESTAMP(3),
duration DOUBLE, — 停留时长
page_url STRING,
device_type STRING,
WATERMARK FOR behavior_time AS behavior_time – INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
t_env.execute_sql(source_ddl)
# 2. 实时特征计算(滑动窗口聚合)
feature_query = """
CREATE VIEW user_behavior_features AS
SELECT
user_id,
HOP_START(behavior_time, INTERVAL '1' MINUTE, INTERVAL '10' MINUTE) AS window_start,
HOP_END(behavior_time, INTERVAL '1' MINUTE, INTERVAL '10' MINUTE) AS window_end,
— 基础统计特征
COUNT(*) AS total_actions,
SUM(CASE WHEN behavior_type = 'click' THEN 1 ELSE 0 END) AS click_count,
SUM(CASE WHEN behavior_type = 'purchase' THEN 1 ELSE 0 END) AS purchase_count,
AVG(duration) AS avg_duration,
— 时间窗口特征
MAX(behavior_time) AS last_action_time,
COUNT(DISTINCT item_id) AS unique_items,
— 设备特征
LAST_VALUE(device_type) AS latest_device,
— 实时比率特征
SUM(CASE WHEN behavior_type = 'purchase' THEN 1 ELSE 0 END) * 1.0 /
GREATEST(COUNT(*), 1) AS conversion_rate
FROM user_behavior_stream
GROUP BY
HOP(behavior_time, INTERVAL '1' MINUTE, INTERVAL '10' MINUTE),
user_id
"""
t_env.execute_sql(feature_query)
return t_env.sql_query("SELECT * FROM user_behavior_features")
# 3. 特征标准化管道
def create_feature_pipeline(feature_table):
"""创建特征预处理管道"""
# 字符串特征索引化
device_indexer = StringIndexer() \\
.set_input_col("latest_device") \\
.set_output_col("device_index")
# 数值特征标准化
numerical_features = ['total_actions', 'click_count', 'purchase_count',
'avg_duration', 'unique_items', 'conversion_rate']
# 特征向量组装
assembler = VectorAssembler() \\
.set_input_cols(numerical_features + ['device_index']) \\
.set_output_col("features")
# 特征标准化
scaler = StandardScaler() \\
.set_input_col("features") \\
.set_output_col("scaled_features")
# 构建特征管道
pipeline = Pipeline(stages=[device_indexer, assembler, scaler])
return pipeline
# 执行特征工程
feature_table = create_user_behavior_features()
feature_pipeline = create_feature_pipeline(feature_table)
# 转换特征
model = feature_pipeline.fit(feature_table)
transformed_features = model.transform(feature_table)
# 注册特征表供后续使用
t_env.create_temporary_view("scaled_features", transformed_features)
2.2 时间序列特征
实时时间序列特征提取。
from pyflink.table.udf import udf
from pyflink.table import DataTypes
from pyflink.table.expressions import col, call
import numpy as np
from scipy import stats
# 自定义时间序列特征UDF
@udf(result_type=DataTypes.STRING(), func_type="pandas")
def extract_time_series_features(values, timestamps):
"""提取时间序列统计特征"""
if len(values) == 0:
return "{}"
series = pd.Series(values, index=pd.to_datetime(timestamps))
features = {
# 统计特征
'mean': np.mean(values),
'std': np.std(values),
'skew': stats.skew(values),
'kurtosis': stats.kurtosis(values),
# 时间特征
'trend': np.polyfit(range(len(values)), values, 1)[0], # 趋势斜率
'volatility': np.std(np.diff(values)) / np.mean(np.abs(np.diff(values))) if len(values) > 1 else 0,
# 窗口特征
'max_value': np.max(values),
'min_value': np.min(values),
'range': np.ptp(values),
# 百分位数
'q25': np.percentile(values, 25),
'q75': np.percentile(values, 75)
}
return str(features)
# 实时股票数据特征提取
def create_stock_features_pipeline():
"""股票实时特征工程"""
stock_source_ddl = """
CREATE TABLE stock_tick_data (
symbol STRING,
price DOUBLE,
volume BIGINT,
timestamp TIMESTAMP(3),
WATERMARK FOR timestamp AS timestamp – INTERVAL '1' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'stock_ticks',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
t_env.execute_sql(stock_source_ddl)
# 滑动窗口特征计算
feature_query = """
CREATE VIEW stock_technical_features AS
SELECT
symbol,
TUMBLE_START(timestamp, INTERVAL '5' MINUTE) AS window_start,
TUMBLE_END(timestamp, INTERVAL '5' MINUTE) AS window_end,
— 价格特征
FIRST_VALUE(price) AS open_price,
LAST_VALUE(price) AS close_price,
MAX(price) AS high_price,
MIN(price) AS low_price,
AVG(price) AS avg_price,
— 成交量特征
SUM(volume) AS total_volume,
AVG(volume) AS avg_volume,
— 技术指标
(LAST_VALUE(price) – FIRST_VALUE(price)) / FIRST_VALUE(price) AS return_5min,
(MAX(price) – MIN(price)) / FIRST_VALUE(price) AS volatility_5min,
— 时间序列特征(自定义UDF)
extract_time_series_features(
COLLECT(price) OVER w,
COLLECT(timestamp) OVER w
) AS time_series_features
FROM stock_tick_data
GROUP BY
TUMBLE(timestamp, INTERVAL '5' MINUTE),
symbol
WINDOW w AS (
PARTITION BY symbol
ORDER BY timestamp
ROWS BETWEEN 29 PRECEDING AND CURRENT ROW — 30个数据点窗口
)
"""
t_env.execute_sql(feature_query)
return t_env.sql_query("SELECT * FROM stock_technical_features")
# 注册UDF
t_env.create_temporary_system_function("extract_time_series_features", extract_time_series_features)
stock_features = create_stock_features_pipeline()
3. 在线学习与模型更新
3.1 增量学习实现
流式数据上的模型持续学习。
from pyflink.ml.classification import LogisticRegression
from pyflink.ml.evaluation import BinaryClassificationEvaluator
from pyflink.ml.tuning import CrossValidator, ParamGrid
from pyflink.table import Table
import json
class StreamingModelManager:
"""流式模型管理器"""
def __init__(self, t_env):
self.t_env = t_env
self.current_model = None
self.model_version = 0
self.evaluation_results = []
def create_training_stream(self):
"""创建带标签的训练数据流"""
training_ddl = """
CREATE TABLE training_data_stream (
features ARRAY<DOUBLE>,
label INT, — 0/1 二分类标签
event_time TIMESTAMP(3),
data_source STRING,
WATERMARK FOR event_time AS event_time – INTERVAL '10' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'model_training',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
self.t_env.execute_sql(training_ddl)
return self.t_env.sql_query("SELECT * FROM training_data_stream")
def incremental_training_pipeline(self, training_data: Table, model_update_interval='10 minutes'):
"""增量训练管道"""
# 创建逻辑回归模型
lr = LogisticRegression() \\
.set_features_col('features') \\
.set_label_col('label') \\
.set_prediction_col('prediction') \\
.set_max_iter(100) \\
.set_reg_param(0.01)
# 时间窗口批训练
training_query = f"""
CREATE VIEW batched_training_data AS
SELECT
TUMBLE_START(event_time, INTERVAL '{model_update_interval}') AS batch_start,
TUMBLE_END(event_time, INTERVAL '{model_update_interval}') AS batch_end,
COLLECT(features) AS features_batch,
COLLECT(label) AS labels_batch,
COUNT(*) AS batch_size
FROM training_data
GROUP BY TUMBLE(event_time, INTERVAL '{model_update_interval}')
"""
self.t_env.execute_sql(training_query)
batched_data = self.t_env.sql_query("SELECT * FROM batched_training_data")
return lr, batched_data
def model_evaluation(self, predictions: Table):
"""模型性能评估"""
evaluator = BinaryClassificationEvaluator() \\
.set_label_col('label') \\
.set_prediction_col('prediction') \\
.set_metric_name('areaUnderROC')
# 计算AUC
auc = evaluator.transform(predictions)
# 评估结果写入监控系统
eval_query = """
INSERT INTO model_performance_metrics
SELECT
CURRENT_TIMESTAMP AS eval_time,
'logistic_regression' AS model_type,
{model_version} AS model_version,
auc AS auc_score,
batch_size AS sample_size
FROM evaluation_results
"""
return auc
# 使用示例
model_manager = StreamingModelManager(t_env)
training_data = model_manager.create_training_stream()
lr_model, batched_data = model_manager.incremental_training_pipeline(training_data)
# 模型拟合
model = lr_model.fit(batched_data)
# 实时预测
predictions = model.transform(training_data)
# 性能评估
auc_score = model_manager.model_evaluation(predictions)
3.2 模型版本管理与A/B测试
生产环境模型部署策略。
from pyflink.ml import Model
from pyflink.table import Table, TableResult
import pickle
import hashlib
from datetime import datetime
class ModelVersionManager:
"""模型版本管理"""
def __init__(self, model_storage_path="hdfs:///models/"):
self.model_storage_path = model_storage_path
self.deployed_models = {} # model_name -> version_info
def save_model(self, model: Model, model_name: str, metadata: dict) –> str:
"""保存模型到存储系统"""
# 生成版本号(基于时间戳和内容哈希)
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
model_bytes = pickle.dumps(model)
model_hash = hashlib.md5(model_bytes).hexdigest()[:8]
version_id = f"{timestamp}_{model_hash}"
# 模型文件路径
model_path = f"{self.model_storage_path}{model_name}/{version_id}.pkl"
metadata_path = f"{self.model_storage_path}{model_name}/{version_id}.json"
# 保存模型和元数据
with open(model_path, 'wb') as f:
pickle.dump(model, f)
metadata.update({
'version_id': version_id,
'created_at': timestamp,
'model_path': model_path,
'model_size': len(model_bytes)
})
with open(metadata_path, 'w') as f:
json.dump(metadata, f, indent=2)
return version_id
def load_model(self, model_name: str, version_id: str = None) –> Model:
"""加载特定版本模型"""
if version_id is None:
# 加载最新版本
version_id = self.get_latest_version(model_name)
model_path = f"{self.model_storage_path}{model_name}/{version_id}.pkl"
with open(model_path, 'rb') as f:
model = pickle.load(f)
return model
def deploy_model(self, model_name: str, version_id: str, traffic_ratio: float = 1.0):
"""部署模型到生产环境"""
if model_name not in self.deployed_models:
self.deployed_models[model_name] = []
deployment_info = {
'version_id': version_id,
'deployed_at': datetime.now().isoformat(),
'traffic_ratio': traffic_ratio,
'is_active': True
}
self.deployed_models[model_name].append(deployment_info)
# 更新模型路由表
self._update_model_routing_table(model_name, version_id, traffic_ratio)
def _update_model_routing_table(self, model_name: str, version_id: str, traffic_ratio: float):
"""更新模型路由配置"""
routing_ddl = f"""
CREATE TABLE IF NOT EXISTS model_routing_config (
model_name STRING,
version_id STRING,
traffic_ratio DOUBLE,
is_active BOOLEAN,
updated_at TIMESTAMP(3)
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/ml_models',
'table-name' = 'model_routing'
)
"""
self.t_env.execute_sql(routing_ddl)
# 更新路由配置
update_sql = f"""
INSERT INTO model_routing_config
VALUES (
'{model_name}',
'{version_id}',
{traffic_ratio},
true,
CURRENT_TIMESTAMP
)
"""
self.t_env.execute_sql(update_sql)
# A/B测试管道
def create_ab_testing_pipeline(t_env, model_a: Model, model_b: Model, split_ratio: float = 0.5):
"""创建A/B测试管道"""
# 创建预测数据流
prediction_ddl = """
CREATE TABLE prediction_requests (
request_id STRING,
features ARRAY<DOUBLE>,
user_id BIGINT,
request_time TIMESTAMP(3),
WATERMARK FOR request_time AS request_time – INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'prediction_requests',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
t_env.execute_sql(prediction_ddl)
requests = t_env.sql_query("SELECT * FROM prediction_requests")
# 随机分流(A/B测试)
ab_split_query = """
CREATE VIEW ab_test_assignments AS
SELECT
request_id,
features,
user_id,
request_time,
CASE
WHEN MOD(ABS(HASH(user_id)), 100) < {split_ratio * 100} THEN 'A'
ELSE 'B'
END AS experiment_group
FROM prediction_requests
""".format(split_ratio=split_ratio * 100)
t_env.execute_sql(ab_split_query)
assignments = t_env.sql_query("SELECT * FROM ab_test_assignments")
# 分别应用两个模型
predictions_a = model_a.transform(
assignments.filter("experiment_group = 'A'")
).select("request_id, prediction as prediction_a, request_time")
predictions_b = model_b.transform(
assignments.filter("experiment_group = 'B'")
).select("request_id, prediction as prediction_b, request_time")
# 合并预测结果
combined_predictions = predictions_a.union_all(predictions_b)
return combined_predictions
4. 实时异常检测
4.1 流式异常检测算法
基于统计和机器学习的实时异常检测。
from pyflink.ml.clustering import KMeans
from pyflink.ml.feature import StandardScaler
from pyflink.table.udf import udf, ScalarFunction
from pyflink.table import DataTypes
import numpy as np
from sklearn.covariance import EllipticEnvelope
from pyflink.common import Row
class StreamingAnomalyDetector:
"""流式异常检测器"""
def __init__(self, t_env):
self.t_env = t_env
self.clustering_model = None
self.scaler = StandardScaler()
def create_metrics_stream(self):
"""创建系统指标监控流"""
metrics_ddl = """
CREATE TABLE system_metrics (
hostname STRING,
metric_name STRING,
metric_value DOUBLE,
timestamp TIMESTAMP(3),
tags MAP<STRING, STRING>,
WATERMARK FOR timestamp AS timestamp – INTERVAL '30' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'system_metrics',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
self.t_env.execute_sql(metrics_ddl)
return self.t_env.sql_query("SELECT * FROM system_metrics")
def extract_time_window_features(self, metrics_table: Table, window_size='5 minutes'):
"""提取时间窗口特征"""
feature_query = f"""
CREATE VIEW metric_window_features AS
SELECT
hostname,
metric_name,
HOP_START(timestamp, INTERVAL '1' MINUTE, INTERVAL '{window_size}') AS window_start,
HOP_END(timestamp, INTERVAL '1' MINUTE, INTERVAL '{window_size}') AS window_end,
— 统计特征
AVG(metric_value) AS avg_value,
STDDEV(metric_value) AS std_value,
MAX(metric_value) AS max_value,
MIN(metric_value) AS min_value,
(MAX(metric_value) – MIN(metric_value)) AS range_value,
— 变化率特征
(LAST_VALUE(metric_value) – FIRST_VALUE(metric_value)) /
NULLIF(FIRST_VALUE(metric_value), 0) AS change_rate,
— 百分位数特征(近似)
APPROX_PERCENTILE(metric_value, 0.95) AS p95_value,
APPROX_PERCENTILE(metric_value, 0.99) AS p99_value,
— 异常分数(基于历史分布)
ABS(metric_value – AVG(metric_value)) / NULLIF(STDDEV(metric_value), 0) AS z_score
FROM system_metrics
GROUP BY
HOP(timestamp, INTERVAL '1' MINUTE, INTERVAL '{window_size}'),
hostname, metric_name
"""
self.t_env.execute_sql(feature_query)
return self.t_env.sql_query("SELECT * FROM metric_window_features")
def clustering_based_anomaly_detection(self, features_table: Table):
"""基于聚类的异常检测"""
# 特征向量组装
feature_cols = ['avg_value', 'std_value', 'max_value', 'min_value',
'range_value', 'change_rate', 'p95_value', 'p99_value', 'z_score']
from pyflink.ml.feature import VectorAssembler
assembler = VectorAssembler() \\
.set_input_cols(feature_cols) \\
.set_output_col("features")
# 特征标准化
scaler = StandardScaler() \\
.set_input_col("features") \\
.set_output_col("scaled_features")
# KMeans聚类
kmeans = KMeans() \\
.set_k(3) \\ # 3个聚类:正常、警告、异常
.set_features_col("scaled_features") \\
.set_prediction_col("cluster")
# 构建管道
from pyflink.ml import Pipeline
pipeline = Pipeline(stages=[assembler, scaler, kmeans])
# 训练聚类模型
model = pipeline.fit(features_table)
clustered_data = model.transform(features_table)
return clustered_data
@udf(result_type=DataTypes.DOUBLE())
def calculate_anomaly_score(features, cluster_center, historical_centers):
"""计算异常分数"""
if not features or not cluster_center:
return 0.0
current_point = np.array(features)
center = np.array(cluster_center)
# 计算到聚类中心的距离
distance = np.linalg.norm(current_point – center)
# 计算到历史中心的平均距离(动态阈值)
if historical_centers and len(historical_centers) > 0:
historical_distances = [
np.linalg.norm(current_point – np.array(center))
for center in historical_centers
]
avg_historical_distance = np.mean(historical_distances)
anomaly_score = distance / avg_historical_distance
else:
anomaly_score = distance
return float(anomaly_score)
# 使用异常检测器
detector = StreamingAnomalyDetector(t_env)
metrics_stream = detector.create_metrics_stream()
features = detector.extract_time_window_features(metrics_stream)
anomalies = detector.clustering_based_anomaly_detection(features)
# 注册UDF
t_env.create_temporary_system_function("calculate_anomaly_score",
detector.calculate_anomaly_score)
# 计算异常分数
anomaly_scored = anomalies.add_columns(
call('calculate_anomaly_score',
anomalies.scaled_features,
anomalies.cluster_center,
anomalies.historical_centers).alias('anomaly_score')
)
# 定义异常阈值
anomaly_threshold = 2.0 # 2倍平均距离
anomalies_detected = anomaly_scored.filter(f"anomaly_score > {anomaly_threshold}")
# 告警输出
alert_ddl = """
CREATE TABLE anomaly_alerts (
hostname STRING,
metric_name STRING,
anomaly_score DOUBLE,
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
alert_time TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'topic' = 'anomaly_alerts',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
t_env.execute_sql(alert_ddl)
# 插入告警
anomalies_detected.execute_insert("anomaly_alerts")
5. 模型服务与部署
5.1 实时模型服务
低延迟模型预测服务。
from pyflink.common import Configuration
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
import threading
import time
from concurrent.futures import ThreadPoolExecutor
class ModelService:
"""模型预测服务"""
def __init__(self, t_env, model_endpoint="localhost:8080"):
self.t_env = t_env
self.model_endpoint = model_endpoint
self.model_cache = {} # model_id -> (model, version, loaded_time)
self.cache_ttl = 3600 # 1小时缓存
self.executor = ThreadPoolExecutor(max_workers=10)
def create_prediction_service(self):
"""创建预测服务管道"""
# 预测请求流
request_ddl = """
CREATE TABLE prediction_requests (
request_id STRING,
model_id STRING,
features ARRAY<DOUBLE>,
request_time TIMESTAMP(3),
priority INT, — 优先级:1-高,2-中,3-低
WATERMARK FOR request_time AS request_time – INTERVAL '1' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'prediction_requests',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
self.t_env.execute_sql(request_ddl)
# 预测结果流
result_ddl = """
CREATE TABLE prediction_results (
request_id STRING,
model_id STRING,
prediction DOUBLE,
confidence DOUBLE,
processing_time_ms BIGINT,
model_version STRING,
error_message STRING
) WITH (
'connector' = 'kafka',
'topic' = 'prediction_results',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
)
"""
self.t_env.execute_sql(result_ddl)
@udf(result_type=DataTypes.ROW([
DataTypes.FIELD("prediction", DataTypes.DOUBLE()),
DataTypes.FIELD("confidence", DataTypes.DOUBLE()),
DataTypes.FIELD("processing_time", DataTypes.BIGINT())
]))
def predict_batch(features_batch, model_id):
"""批量预测UDF"""
start_time = time.time()
try:
# 从缓存获取模型
model = self._get_model_from_cache(model_id)
# 批量预测
predictions = model.predict(features_batch)
confidences = model.predict_proba(features_batch)
processing_time = int((time.time() – start_time) * 1000) # 毫秒
return Row(
prediction=float(predictions[0]),
confidence=float(np.max(confidences[0])),
processing_time=processing_time
)
except Exception as e:
return Row(
prediction=–1.0,
confidence=0.0,
processing_time=–1,
error_message=str(e)
)
def _get_model_from_cache(self, model_id: str):
"""从缓存获取模型(带自动刷新)"""
if model_id in self.model_cache:
model, version, loaded_time = self.model_cache[model_id]
# 检查缓存是否过期
if time.time() – loaded_time < self.cache_ttl:
return model
# 从模型仓库加载最新版本
latest_model = self._load_model_from_registry(model_id)
self.model_cache[model_id] = (latest_model, "latest", time.time())
return latest_model
def _load_model_from_registry(self, model_id: str):
"""从模型注册中心加载模型"""
# 这里可以实现从HDFS、S3或模型服务器加载
model_path = f"hdfs:///models/{model_id}/latest.pkl"
# 模拟模型加载
# 实际实现中会从存储系统加载序列化的模型
return f"loaded_model_{model_id}"
def create_low_latency_pipeline(self):
"""创建低延迟预测管道"""
# 优先级处理:高优先级请求优先处理
priority_query = """
CREATE VIEW prioritized_requests AS
SELECT
request_id,
model_id,
features,
request_time,
priority,
— 为不同优先级设置不同超时
CASE
WHEN priority = 1 THEN request_time + INTERVAL '100' MILLISECOND
WHEN priority = 2 THEN request_time + INTERVAL '500' MILLISECOND
WHEN priority = 3 THEN request_time + INTERVAL '1000' MILLISECOND
END AS timeout_time
FROM prediction_requests
"""
self.t_env.execute_sql(priority_query)
# 注册预测UDF
self.t_env.create_temporary_system_function("predict_batch", self.predict_batch)
# 实时预测
prediction_query = """
CREATE VIEW realtime_predictions AS
SELECT
request_id,
model_id,
prediction,
confidence,
processing_time_ms,
'v1.0' AS model_version,
error_message
FROM prioritized_requests,
LATERAL TABLE(predict_batch(features, model_id)) AS T(
prediction DOUBLE,
confidence DOUBLE,
processing_time BIGINT
)
WHERE CURRENT_TIMESTAMP < timeout_time — 超时过滤
"""
self.t_env.execute_sql(prediction_query)
# 写入预测结果
t_env.sql_query("SELECT * FROM realtime_predictions") \\
.execute_insert("prediction_results")
# 启动预测服务
model_service = ModelService(t_env)
model_service.create_prediction_service()
model_service.create_low_latency_pipeline()
6. ML Pipeline监控与运维
6.1 全链路监控体系
机器学习管道的可观测性。
from pyflink.ml import Pipeline
from pyflink.table import Table
import json
import time
from datetime import datetime
class MLOpsMonitor:
"""MLOps监控系统"""
def __init__(self, t_env):
self.t_env = t_env
self.metrics = {}
def create_monitoring_tables(self):
"""创建监控数据表"""
# 模型性能监控
performance_ddl = """
CREATE TABLE model_performance_metrics (
model_id STRING,
metric_name STRING,
metric_value DOUBLE,
timestamp TIMESTAMP(3),
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
sample_size BIGINT
) WITH (
'connector' = 'elasticsearch',
'hosts' = 'http://elasticsearch:9200',
'index' = 'model-metrics-{now/d}'
)
"""
self.t_env.execute_sql(performance_ddl)
# 数据质量监控
data_quality_ddl = """
CREATE TABLE data_quality_metrics (
pipeline_id STRING,
check_point STRING,
check_type STRING,
actual_value DOUBLE,
expected_value DOUBLE,
status STRING,
check_time TIMESTAMP(3)
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/ml_monitoring',
'table-name' = 'data_quality_checks'
)
"""
self.t_env.execute_sql(data_quality_ddl)
# 系统资源监控
resource_ddl = """
CREATE TABLE resource_metrics (
worker_id STRING,
cpu_usage DOUBLE,
memory_usage DOUBLE,
gc_time_ms BIGINT,
record_throughput DOUBLE,
timestamp TIMESTAMP(3)
) WITH (
'connector' = 'prometheus',
'url' = 'http://prometheus:9090'
)
"""
self.t_env.execute_sql(resource_ddl)
def monitor_feature_drift(self, feature_table: Table, reference_stats: dict):
"""监控特征漂移"""
drift_query = f"""
CREATE VIEW feature_drift_monitor AS
SELECT
CURRENT_TIMESTAMP AS check_time,
'feature_drift' AS check_type,
feature_name,
current_stats->'mean' AS current_mean,
{reference_stats['mean']} AS reference_mean,
ABS(current_stats->'mean' – {reference_stats['mean']}) /
NULLIF({reference_stats['std']}, 0) AS drift_score,
CASE
WHEN ABS(current_stats->'mean' – {reference_stats['mean']}) /
NULLIF({reference_stats['std']}, 0) > 3 THEN 'HIGH_DRIFT'
WHEN ABS(current_stats->'mean' – {reference_stats['mean']}) /
NULLIF({reference_stats['std']}, 0) > 2 THEN 'MEDIUM_DRIFT'
ELSE 'LOW_DRIFT'
END AS drift_level
FROM (
SELECT
feature_name,
MAP['mean', AVG(value), 'std', STDDEV(value)] AS current_stats
FROM feature_table
CROSS JOIN UNNEST(features) AS T(feature_name, value)
GROUP BY feature_name
)
"""
self.t_env.execute_sql(drift_query)
return self.t_env.sql_query("SELECT * FROM feature_drift_monitor")
def monitor_concept_drift(self, predictions: Table, window_size='1 hour'):
"""监控概念漂移"""
drift_detection_query = f"""
CREATE VIEW concept_drift_monitor AS
SELECT
TUMBLE_START(event_time, INTERVAL '{window_size}') AS window_start,
TUMBLE_END(event_time, INTERVAL '{window_size}') AS window_end,
model_id,
— 准确率变化
AVG(CASE WHEN actual = predicted THEN 1.0 ELSE 0.0 END) AS accuracy,
— 准确率滑动窗口统计
AVG(accuracy) OVER (
PARTITION BY model_id
ORDER BY window_start
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
) AS accuracy_ma_7,
— 准确率变化率
(accuracy – LAG(accuracy, 1) OVER w) /
NULLIF(LAG(accuracy, 1) OVER w, 0) AS accuracy_change_rate,
— 漂移检测
CASE
WHEN accuracy < 0.8 THEN 'HIGH_DRIFT'
WHEN accuracy_change_rate < -0.1 THEN 'MEDIUM_DRIFT'
ELSE 'STABLE'
END AS drift_status
FROM predictions
GROUP BY
TUMBLE(event_time, INTERVAL '{window_size}'),
model_id
WINDOW w AS (PARTITION BY model_id ORDER BY window_start)
"""
self.t_env.execute_sql(drift_detection_query)
return self.t_env.sql_query("SELECT * FROM concept_drift_monitor")
def create_alert_system(self):
"""创建告警系统"""
alert_rules_ddl = """
CREATE TABLE alert_rules (
rule_id STRING,
rule_name STRING,
rule_condition STRING,
severity STRING, — CRITICAL, HIGH, MEDIUM, LOW
notification_channel STRING,
cooldown_minutes INT
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/ml_monitoring',
'table-name' = 'alert_rules'
)
"""
self.t_env.execute_sql(alert_rules_ddl)
# 告警规则示例
alert_rules = [
('DRIFT_01', 'high_feature_drift', 'drift_score > 3.0', 'HIGH', 'slack', 30),
('PERF_01', 'low_accuracy', 'accuracy < 0.7', 'CRITICAL', 'pagerduty', 5),
('DATA_01', 'missing_data', 'sample_size < 1000', 'MEDIUM', 'email', 60)
]
for rule in alert_rules:
insert_sql = f"""
INSERT INTO alert_rules
VALUES {rule}
"""
self.t_env.execute_sql(insert_sql)
# 初始化监控系统
mlops = MLOpsMonitor(t_env)
mlops.create_monitoring_tables()
mlops.create_alert_system()
# 启动监控作业
feature_drift = mlops.monitor_feature_drift(features_table, reference_stats)
concept_drift = mlops.monitor_concept_drift(predictions)
# 告警触发
alert_trigger_query = """
CREATE VIEW triggered_alerts AS
SELECT
rule.rule_id,
rule.rule_name,
rule.severity,
rule.notification_channel,
metrics.*,
CURRENT_TIMESTAMP AS alert_time
FROM alert_rules rule
JOIN monitoring_metrics metrics ON eval(rule.rule_condition)
WHERE rule.severity IN ('CRITICAL', 'HIGH')
"""
t_env.execute_sql(alert_trigger_query)
7. 总结
PyFlink与ML pipeline的集成为流式机器学习提供了强大的基础设施。关键优势包括:
核心价值
1. 实时特征工程: 滑动窗口聚合、时间序列特征、实时统计
2. 在线学习: 增量训练、模型动态更新、概念漂移检测
3. 生产就绪: 模型版本管理、A/B测试、监控告警
4. 低延迟服务: 批量预测优化、优先级处理、缓存机制
最佳实践
1. 特征设计: 平衡实时性与历史上下文,使用滑动窗口聚合
2. 模型更新: 定期重训练与增量学习结合,监控模型性能衰减
3. 资源优化: 合理设置并行度,使用缓存减少重复计算
4. 监控完备: 全链路可观测性,数据漂移与概念漂移检测
适用场景
- 实时推荐系统
- 金融风控与欺诈检测
- IoT设备异常预测
- 实时用户行为分析
- 系统监控与智能运维
这种架构使得机器学习模型能够真正实时响应数据变化,在动态业务环境中保持高预测准确性。

