
Granite TimeSeries FlowState R1数据管道构建从CSV到预测的全流程自动化你是不是也遇到过这样的麻烦手头有一堆CSV格式的销售数据、设备传感器日志或者网站流量记录想用Granite TimeSeries FlowState R1这个时间序列模型来做预测但每次都得手动下载数据、写脚本清洗、再调用模型最后还得把结果存起来。整个过程繁琐又容易出错更别提要每天、每周定时跑一遍了。今天我就带你从零开始搭建一个全自动化的数据管道。这个管道能帮你定时抓取数据、自动清洗、调用模型预测最后把结果稳稳当当地存进数据库。整个过程就像设置了一个智能流水线你只需要配置一次它就能在后台默默为你工作把预测结果准时送到你手上。我们主要会用Python来写核心的处理逻辑再用一个叫Prefect的轻量级工具来负责定时调度和任务监控。整个流程清晰简单即使你之前没怎么接触过数据管道跟着步骤走也能轻松搞定。1. 先理清思路我们的自动化管道要做什么在动手写代码之前我们得先想清楚这个管道到底要完成哪些事。你可以把它想象成一个制作预测报告的智能工厂每个车间都有明确的任务。数据来源车间我们的原材料是CSV文件。它们可能放在公司的FTP服务器上也可能在云存储比如阿里云OSS、AWS S3里。管道的第一个任务就是定时去这些地方把最新的数据文件“搬”回来。数据清洗车间搬回来的原材料原始CSV可能有点脏乱。比如有些行数据缺失了日期格式不统一或者有些列我们根本用不上。这个车间负责把数据整理干净转换成模型喜欢吃的“标准餐”。模型预测车间这是核心环节。我们把清洗好的数据按照Granite TimeSeries FlowState R1模型API要求的格式打包好发送过去。然后等待模型“消化”数据给我们吐出一份未来的预测结果。结果交付车间预测结果不能只放在程序内存里。我们需要把它持久化地存起来比如写入MySQL、PostgreSQL这类数据库方便后续的报表系统查询或者直接生成一个Excel、PDF报告发送给相关同事。调度控制中心上面四个车间需要在正确的时间、按照正确的顺序工作。比如必须等数据下载并清洗完成后才能发起预测。Prefect就是这个控制中心它负责定时触发整个流程并确保一个环节出错时能及时通知我们。理清了这五个部分我们心里就有了一张清晰的蓝图。接下来我们就从最基础的环节开始搭建每一个“车间”。2. 搭建你的Python数据处理环境工欲善其事必先利其器。我们先来把Python环境准备好安装好需要的“工具包”。我建议你使用conda或者venv来创建一个独立的Python环境这样不会和你电脑上其他的项目冲突。这里我用conda举例# 创建一个名为 granite-pipeline 的新环境指定Python版本 conda create -n granite-pipeline python3.9 # 激活这个环境 conda activate granite-pipeline环境激活后我们来安装这个项目最核心的几个Python库。你可以把下面这些命令一次性复制执行。pip install pandas numpy requests prefect prefect-sqlalchemy我来简单说说这几个库是干嘛的pandas处理表格数据CSV的神器清洗、转换都靠它。numpypandas的好搭档进行高效的数值计算。requests用来发送HTTP请求我们靠它来调用Granite模型的API。prefect今天的主角之一负责工作流的编排、调度和监控。prefect-sqlalchemyPrefect的一个集成包能让我们更方便地把数据写入各种数据库。如果你的数据源是FTP可能还需要ftplibPython自带或paramiko用于SFTP如果数据在云存储则需要对应的SDK比如boto3AWS S3或oss2阿里云OSS。这里我们先以最通用的本地文件或HTTP链接为例你可以根据实际情况调整。安装完成后新建一个项目文件夹比如granite_data_pipeline我们所有的代码文件都会放在这里。3. 编写核心数据处理模块现在我们来逐一实现蓝图里的各个“车间”。我们会把每个功能都写成一个独立的Python函数这样逻辑清晰也方便测试和复用。3.1 从各种来源获取数据数据可能来自不同地方。我们先写一个通用的数据加载函数它能根据文件路径或URL来读取CSV。import pandas as pd import os def load_data_from_source(source_path): 从指定路径加载CSV数据。 支持本地文件路径和HTTP/HTTPS URL。 参数: source_path (str): 数据源路径例如 /data/raw.csv 或 https://example.com/data.csv 返回: pandas.DataFrame: 加载后的数据框 try: # 判断是否是网络路径 if source_path.startswith((http://, https://)): print(f正在从网络地址加载数据: {source_path}) # 这里简单使用pandas直接读取对于复杂情况可能需要添加headers等参数 df pd.read_csv(source_path) else: # 本地文件路径 print(f正在从本地加载数据: {source_path}) if not os.path.exists(source_path): raise FileNotFoundError(f文件未找到: {source_path}) df pd.read_csv(source_path) print(f数据加载成功共 {len(df)} 行{len(df.columns)} 列。) return df except Exception as e: # 在实际管道中这里应该记录更详细的日志并可能触发重试或告警 print(f加载数据时发生错误: {e}) # 可以选择将错误向上抛出由工作流引擎处理 raise # 示例用法 # df_raw load_data_from_source(./data/sales_data.csv) # 或者 # df_raw load_data_from_source(https://your-data-server.com/monthly_traffic.csv)3.2 把数据清洗干净原始数据往往不太规整。清洗函数的目标是产出模型能直接“食用”的干净数据。def clean_and_transform_data(df_raw): 对原始数据进行清洗和转换。 参数: df_raw (pandas.DataFrame): 原始数据框 返回: pandas.DataFrame: 清洗转换后的数据框 df df_raw.copy() # 避免修改原始数据 print(开始数据清洗与转换...) # 1. 处理日期时间列 (时间序列分析的核心) # 假设你的数据里有一个叫 timestamp 的列 if timestamp in df.columns: df[timestamp] pd.to_datetime(df[timestamp], errorscoerce) # 将日期设为主索引这对时间序列分析很重要 df.set_index(timestamp, inplaceTrue) print(已将 timestamp 列转换为日期时间索引。) else: print(警告未找到 timestamp 列请确保数据包含时间信息。) # 2. 处理缺失值 # 对于数值列用前后值的平均值填充根据业务逻辑选择合适方法 numeric_cols df.select_dtypes(include[number]).columns if not numeric_cols.empty: df[numeric_cols] df[numeric_cols].interpolate(methodlinear).fillna(methodbfill).fillna(methodffill) print(f已对数值列 ({list(numeric_cols)}) 进行缺失值插补。) # 3. 移除完全空白的行或列 df.dropna(howall, inplaceTrue) df.dropna(axis1, howall, inplaceTrue) # 4. 重命名列使其更规范可选 # df.rename(columns{old_name: new_name}, inplaceTrue) # 5. 数据采样或聚合如果需要 # 例如将高频数据聚合为每日数据 # df_resampled df.resample(D).mean() print(f清洗完成。数据形状: {df.shape}) return df # 或者返回 df_resampled3.3 调用模型API获取预测这是最令人期待的一步。你需要准备好Granite TimeSeries FlowState R1模型的API端点URL和访问凭证如API Key。import requests import json import time def call_granite_api_for_prediction(clean_data_df, api_endpoint, api_key, forecast_horizon7): 将清洗后的数据发送到Granite模型API获取预测结果。 参数: clean_data_df (pandas.DataFrame): 清洗后的数据框应包含时间索引和数值列。 api_endpoint (str): Granite模型API的完整URL。 api_key (str): 用于认证的API密钥。 forecast_horizon (int): 需要预测的未来步长例如预测未来7天。 返回: dict: API返回的预测结果通常是JSON格式。 print(f准备调用Granite模型API: {api_endpoint}) # 1. 将DataFrame转换为API要求的格式 # 这里需要根据Granite API的实际文档来构造请求体 # 假设API要求一个包含timestamps和values列表的JSON # 我们取第一列数值数据作为示例 if clean_data_df.shape[1] 0: raise ValueError(清洗后的数据没有数值列可供预测。) # 假设我们使用第一列数值数据进行预测 target_series_name clean_data_df.columns[0] series_data clean_data_df[target_series_name].dropna().tolist() timestamps clean_data_df.index.astype(str).tolist()[-len(series_data):] # 确保时间戳和值长度一致 request_payload { series: { timestamps: timestamps, values: series_data }, forecast_horizon: forecast_horizon, # 可能还有其他参数如置信区间、模型配置等 } # 2. 设置请求头通常包含认证信息和内容类型 headers { Authorization: fBearer {api_key}, Content-Type: application/json } # 3. 发送POST请求 try: response requests.post( api_endpoint, datajson.dumps(request_payload), headersheaders, timeout60 # 设置超时时间 ) response.raise_for_status() # 如果状态码不是200抛出异常 # 4. 解析响应 prediction_result response.json() print(fAPI调用成功预测未来 {forecast_horizon} 步。) return prediction_result except requests.exceptions.RequestException as e: print(f调用API时发生网络或请求错误: {e}) if hasattr(e, response) and e.response is not None: print(f错误响应内容: {e.response.text}) raise except json.JSONDecodeError as e: print(f解析API响应JSON时出错: {e}) raise3.4 把预测结果存起来拿到预测结果后我们需要把它持久化。这里展示两种常见方式存入数据库和生成CSV报告。from sqlalchemy import create_engine from prefect_sqlalchemy import SqlAlchemyConnector import pandas as pd def save_results_to_database(prediction_result, table_namemodel_predictions): 将预测结果保存到数据库。 参数: prediction_result (dict): 模型API返回的结果。 table_name (str): 要写入的数据库表名。 # 假设prediction_result里有一个forecast键对应预测值列表 forecast_data prediction_result.get(forecast, []) # 假设还有对应的未来时间戳 forecast_timestamps prediction_result.get(forecast_timestamps, []) if not forecast_data: print(警告预测结果为空跳过数据库写入。) return # 将数据转换为DataFrame df_to_save pd.DataFrame({ forecast_timestamp: forecast_timestamps, predicted_value: forecast_data, created_at: pd.Timestamp.now() # 记录生成时间 }) # 使用Prefect的SqlAlchemyConnector管理数据库连接更安全、可配置 # 首先你需要在Prefect的配置中设置好数据库连接信息如URL print(f正在将预测结果写入数据库表 {table_name}...) try: # 这里演示直接使用sqlalchemy创建引擎生产环境建议用Connector # 请将下面的连接字符串替换为你自己的 # engine create_engine(postgresql://user:passwordlocalhost/dbname) # df_to_save.to_sql(table_name, engine, if_existsappend, indexFalse) # 更佳实践使用Prefect的Block来管理连接需提前在UI或代码中配置 # from prefect.blocks.system import Secret # connection_string Secret.load(my-db-connection-string).get() # engine create_engine(connection_string) # df_to_save.to_sql(table_name, engine, if_existsappend, indexFalse) print(此处为演示实际代码需配置真实的数据库连接) print(f模拟将 {len(df_to_save)} 条预测记录写入 {table_name}。) except Exception as e: print(f写入数据库时发生错误: {e}) raise def save_results_to_csv(prediction_result, output_path./output/predictions.csv): 将预测结果保存为CSV文件。 参数: prediction_result (dict): 模型API返回的结果。 output_path (str): 输出CSV文件的路径。 import os forecast_data prediction_result.get(forecast, []) forecast_timestamps prediction_result.get(forecast_timestamps, []) if not forecast_data: print(警告预测结果为空跳过CSV生成。) return df_to_save pd.DataFrame({ forecast_timestamp: forecast_timestamps, predicted_value: forecast_data }) # 确保输出目录存在 os.makedirs(os.path.dirname(output_path), exist_okTrue) try: df_to_save.to_csv(output_path, indexFalse) print(f预测结果已成功保存至: {output_path}) except Exception as e: print(f保存CSV文件时发生错误: {e}) raise好了四个核心的“车间”模块我们都准备好了。它们各自独立功能明确。接下来我们需要一个“总控室”把它们串联起来并让这个流程能定时自动运行。4. 用Prefect编排自动化工作流Prefect是一个现代的工作流编排工具它比传统的Airflow更轻量、更Pythonic。我们用它将上面的函数组装成一个可靠的任务流水线。4.1 定义Prefect任务和流在Prefect里每个Python函数可以包装成一个task多个task组成一个flow。from prefect import task, flow from prefect.tasks import task_input_hash from datetime import timedelta # 使用task装饰器将我们的函数定义为Prefect任务 # 设置cache_key_fn可以让任务结果在一定条件下被缓存避免重复计算 task(cache_key_fntask_input_hash, cache_expirationtimedelta(hours1)) def task_load_data(source_path): return load_data_from_source(source_path) task def task_clean_data(df_raw): return clean_and_transform_data(df_raw) task def task_get_prediction(clean_df): # 在实际使用中API端点URL和密钥应从安全配置中读取例如Prefect的Secret Block API_ENDPOINT https://api.your-granite-service.com/v1/predict # 替换为真实URL API_KEY your-actual-api-key-here # 替换为真实Key切勿硬编码在代码中 return call_granite_api_for_prediction(clean_df, API_ENDPOINT, API_KEY) task def task_save_to_db(prediction_result): save_results_to_database(prediction_result) return Saved to DB task def task_save_to_csv(prediction_result): output_file f./output/predictions_{pd.Timestamp.now().strftime(%Y%m%d_%H%M)}.csv save_results_to_csv(prediction_result, output_file) return output_file # 使用flow装饰器定义主工作流 flow(namegranite-timeseries-pipeline) def granite_data_pipeline_flow(data_source_path: str ./data/raw_data.csv): 主工作流串联所有任务定义执行逻辑。 print(启动Granite时间序列数据管道...) # 1. 加载数据 raw_data task_load_data(data_source_path) # 2. 清洗数据 clean_data task_clean_data(raw_data) # 3. 调用模型预测 prediction task_get_prediction(clean_data) # 4. 保存结果可以并行执行 db_status task_save_to_db(prediction) csv_file_path task_save_to_csv(prediction) print(f管道执行完毕。数据库状态: {db_status}, CSV文件: {csv_file_path}) return prediction4.2 运行与调度你的工作流现在你可以直接运行这个流来测试整个管道。if __name__ __main__: # 本地运行一次 result granite_data_pipeline_flow(./data/your_data.csv) print(流程运行完成)要让这个流定时自动运行比如每天凌晨1点你不需要写复杂的crontab。Prefect提供了非常简单的调度方式。你只需要在flow装饰器里加上schedule参数即可。from prefect import flow from prefect.client.schemas.schedules import CronSchedule flow( namescheduled-granite-pipeline, scheduleCronSchedule(cron0 1 * * *) # 每天UTC时间1点运行请根据你的时区调整 ) def scheduled_granite_pipeline(): # 这里可以写死数据源路径或从配置中读取 granite_data_pipeline_flow(data_source_pathhttps://your-data-source.com/daily_feed.csv) # 要部署这个带调度的流你需要使用Prefect的服务本地Server或Cloud # 通常在命令行执行: prefect deployment create ./your_script.py:scheduled_granite_pipeline部署到Prefect Server或Cloud后你就可以在漂亮的Web UI上看到任务运行的历史记录、日志、状态成功或失败并设置失败时的告警通知比如发邮件、发Slack消息。5. 让管道更健壮错误处理与日志监控一个生产级的管道不能一出错就崩溃。我们需要给它加上“安全带”和“黑匣子”。1. 任务重试Prefect可以轻松地为任务设置重试机制。比如网络偶尔波动导致API调用失败我们可以让它自动重试几次。task(retries3, retry_delay_seconds10) # 失败后重试3次每次间隔10秒 def task_get_prediction(clean_df): # ... 函数体不变2. 结果缓存对于耗时的、结果不变的任务比如读取某个静态参考数据可以缓存其结果避免每次流运行都重复执行。我们在之前的task_load_data中已经通过cache_key_fn做了演示。3. 完善的日志我们在每个函数里都用print打了日志但在生产环境最好配置更专业的日志库如logging模块将日志输出到文件或日志收集系统方便排查问题。4. 告警通知在Prefect UI中你可以配置“通知策略”。当任务失败、流运行失败时自动发送邮件、Slack消息或Webhook通知到你的团队确保问题能被及时响应。6. 总结走完这一趟一个从数据到预测的自动化管道就初具雏形了。我们先用Python把数据加载、清洗、模型调用、结果存储这几个核心步骤模块化然后用Prefect这个“胶水”把它们粘合起来并赋予了定时调度和状态监控的能力。这套方案的好处是清晰、灵活。每个模块你都可以根据实际需求深入改造比如支持更复杂的数据源、加入更精细的数据验证、处理模型返回的多种结果格式等等。Prefect的生态也在不断丰富你可以很方便地把它和Docker、Kubernetes、各种云服务集成构建更强大、更可靠的数据流水线。下一步你可以尝试把API密钥、数据库连接字符串等敏感信息移到Prefect的Secret管理模块中提升安全性。也可以探索更复杂的流模式比如并行处理多个数据序列或者根据条件决定不同的处理分支。动手试试吧当你看到第一个自动生成的预测报告安静地躺在数据库里时那种解放双手的成就感绝对值得你花时间搭建它。获取更多AI镜像想探索更多AI镜像和应用场景访问 CSDN星图镜像广场提供丰富的预置镜像覆盖大模型推理、图像生成、视频生成、模型微调等多个领域支持一键部署。