Python自动化获取中国气象数据并入库:API调用、数据处理与数据库存储实战

发布时间:2026/7/30 7:39:47
Python自动化获取中国气象数据并入库:API调用、数据处理与数据库存储实战 1. 项目概述与核心价值最近在做一个需要历史天气数据的分析项目第一反应就是去中国气象数据网找找看。但真上手才发现和想象中不太一样网站数据虽然全但一个个文件手动下载、解析、入库工作量巨大且容易出错。于是我花时间研究了一下它的API用Python写了一套自动化的数据获取与入库脚本。今天就把这套从申请API权限、解析数据到自动导入数据库的完整流程和踩过的坑毫无保留地分享出来。无论你是做数据分析、机器学习特征工程还是开发需要气象数据支撑的应用这套方法都能帮你省下大量重复劳动的时间。简单来说这个项目能帮你自动化完成三件事一是通过官方API稳定、合规地获取气象数据二是将获取到的JSON或XML格式的原始数据清洗并转换成结构化的表格数据三是将这些数据高效、准确地存入你指定的数据库比如MySQL、PostgreSQL或SQLite中方便后续的查询与分析。我会附上完整的、可运行的源码你只需要替换几个配置参数就能跑起来。2. 前期准备API申请与环境搭建在写代码之前有两件必须搞定的事情一是去中国气象数据网申请合法的API访问权限二是在本地搭建好Python开发环境。别跳过这一步很多后续的坑都源于前期准备不充分。2.1 中国气象数据网API申请详解中国气象数据网的API服务是需要实名申请和审核的这是为了数据安全和合规使用。整个过程并不复杂但有几个关键点需要注意。首先你需要访问中国气象数据网的官方网站找到“数据服务”或“API接口”相关的板块。通常你需要注册一个账号并完成实名认证。认证通过后在个人中心或开发者平台页面可以找到API服务的申请入口。申请时你需要填写一份申请表说明你的使用目的、用途、预计调用频率等。这里有个小技巧在填写使用目的时尽量具体、非商业、偏向科研或学习例如“用于个人学习Python数据获取与分析技术”或“某高校科研项目的气候数据分析”这样审核通过的概率会高很多审核周期也可能缩短。申请成功后你会获得几个至关重要的凭证API Key有时也叫AppID或Access Key和API Secret或Private Key。此外还会获得一个API请求地址Endpoint。请务必妥善保管这些信息特别是API Secret它相当于你的密码不要直接写在代码里然后上传到公开的代码仓库如GitHub否则可能导致密钥泄露产生不必要的费用或安全风险。正确的做法是使用环境变量或配置文件来管理。注意不同数据产品如地面小时数据、格点预报数据对应的API接口地址和参数可能不同申请时请确认你需要的具体数据产品是否开放了API接口。通常基础的历史气象数据如地面站点的温、压、湿、风、降水都是提供的。2.2 Python开发环境与依赖库安装这个项目对Python版本要求不高Python 3.6及以上都可以。我个人的环境是Python 3.8比较稳定。你需要安装几个核心的库它们各自扮演着关键角色requests: 这是用于发送HTTP请求到气象数据API的库是网络交互的基石。pandas: 数据处理和分析的核心。它能够轻松地将API返回的JSON或字典数据转换为DataFrame并进行清洗、转换。sqlalchemy: 数据库ORM工具。它提供了一个抽象层让我们可以用统一的Python代码来操作不同类型的数据库如MySQL、PostgreSQL、SQLite而无需编写复杂的原生SQL语句。这对于代码的可移植性非常友好。pymysql/psycopg2/sqlite3: 数据库驱动。你需要根据你选择的数据库安装对应的驱动。sqlite3是Python标准库无需安装。你可以使用pip一次性安装这些依赖以MySQL为例pip install requests pandas sqlalchemy pymysql如果你使用PostgreSQL则将pymysql替换为psycopg2-binary。我强烈建议你使用虚拟环境如venv或conda来管理项目依赖这样可以避免不同项目间的库版本冲突。创建一个新的虚拟环境并激活它然后在其中安装上述包是一个好习惯。3. 核心代码模块拆解与原理整个脚本可以清晰地划分为三个功能模块API请求器、数据处理器和数据库操作器。这样的设计遵循了“单一职责原则”每个模块只做一件事使得代码逻辑清晰易于调试和维护。3.1 API请求模块构建稳健的请求链路这个模块的核心任务是构造一个符合气象数据网API规范的HTTP请求并处理可能的网络异常和API错误。我们不能假设网络永远通畅、API永远返回正确结果。首先你需要根据API文档构造请求URL和参数。一个典型的请求可能包含以下参数apiKey: 你的密钥。dataCode: 数据代码例如SURF_CHN_MUL_HOR代表中国地面小时资料。elements: 需要获取的气象要素如TEM温度、PRE气压、RHU相对湿度等多个要素用逗号分隔。stationIDs: 气象站号可以指定单个或多个站。timeRange: 时间范围格式如[202405010000, 202405012300]表示起始和结束时间UTC时间。在代码中我们使用requests库的get或post方法发送请求。这里有一个至关重要的实践必须设置超时timeout参数。我通常设置timeout(10, 30)表示连接超时10秒读取超时30秒。没有超时设置的网络请求在遇到故障时可能会永远挂起。其次必须检查HTTP状态码和API返回的业务状态码。即使HTTP状态码是200成功API也可能因为参数错误、额度不足等原因返回一个包含错误信息的JSON。因此在解析数据之前先判断返回的JSON中是否有returnCode、code或error字段并确认其值为成功如200或success。最后要实施简单的重试机制。对于偶发的网络抖动重试往往能解决问题。我们可以用try-except包裹请求代码在捕获到连接超时、请求异常等错误时进行有限次数的重试比如3次每次重试前等待几秒钟。import requests import time from typing import Dict, Any, Optional class WeatherDataFetcher: def __init__(self, api_key: str, base_url: str): self.api_key api_key self.base_url base_url self.session requests.Session() # 使用Session可以复用TCP连接提升效率 self.session.headers.update({User-Agent: MyWeatherApp/1.0}) # 建议设置User-Agent def fetch_data(self, params: Dict[str, Any], max_retries: int 3) - Optional[Dict]: 获取数据包含重试逻辑 for attempt in range(max_retries): try: # 将API Key加入参数 params[apiKey] self.api_key response self.session.get(self.base_url, paramsparams, timeout(10, 30)) response.raise_for_status() # 如果HTTP状态码不是200将抛出HTTPError异常 data response.json() # 假设API成功返回时有一个returnCode字段且值为200 if data.get(returnCode) ! 200: error_msg data.get(message, Unknown API error) print(fAPI returned an error: {error_msg}) return None # 或者抛出自定义异常 return data except requests.exceptions.Timeout: print(fRequest timeout (attempt {attempt 1}/{max_retries}). Retrying...) except requests.exceptions.ConnectionError: print(fConnection error (attempt {attempt 1}/{max_retries}). Retrying...) except requests.exceptions.RequestException as e: print(fRequest failed: {e}) break # 对于非网络类错误直接跳出重试循环 if attempt max_retries - 1: time.sleep(2 ** attempt) # 指数退避策略等待时间逐渐变长 print(Failed to fetch data after all retries.) return None3.2 数据处理模块从原始JSON到规整DataFrameAPI返回的数据往往不是直接就能用的。它可能是一个嵌套很深的JSON结构包含了元数据如请求信息和实际数据DS字段。我们的目标是把DS里的数据转换成pandas的DataFrame并且让每一列都有清晰的名字。第一步是数据提取。你需要根据API返回的实际结构定位到存放核心数据的字段。在中国气象数据网的许多接口中数据通常位于DS或data字段下它是一个列表列表中的每个元素代表一条记录如一个站点某个时次的数据。第二步是数据扁平化与清洗。原始数据可能包含一些我们不需要的字段或者字段名是编码如V13301我们需要根据数据字典将其映射为可读的名称如temperature。同时要处理缺失值。气象数据中缺失值常用特定值表示比如9999、999999或NaN。我们需要将这些值替换为Python/Pandas能识别的None或np.nan。第三步是类型转换与时间处理。从JSON中解析出来的数字可能是字符串格式需要转换为int或float。时间字段如DATETIME通常是字符串需要利用pandas.to_datetime函数将其转换为Pandas的datetime64类型这对于后续基于时间的筛选和聚合操作至关重要。import pandas as pd import numpy as np class DataProcessor: staticmethod def parse_to_dataframe(raw_data: Dict) - pd.DataFrame: 将API返回的原始数据解析为DataFrame if not raw_data or DS not in raw_data: print(Invalid data format or empty data.) return pd.DataFrame() # 提取数据列表 records raw_data[DS] if not records: return pd.DataFrame() # 创建DataFrame df pd.DataFrame(records) # 示例字段映射根据实际API文档调整 column_mapping { StationID: station_id, DATETIME: obs_time, TEM: temperature, # 温度 PRE: pressure, # 气压 RHU: humidity, # 相对湿度 WIN: wind_speed, # 风速 # ... 其他字段 } df.rename(columnscolumn_mapping, inplaceTrue) # 处理缺失值假设9999代表温度缺失 df.replace({temperature: {9999: np.nan}}, inplaceTrue) # 转换时间列 if obs_time in df.columns: df[obs_time] pd.to_datetime(df[obs_time], format%Y%m%d%H%M, errorscoerce) # 设置时间为索引便于时间序列分析可选 # df.set_index(obs_time, inplaceTrue) # 转换数值列 numeric_cols [temperature, pressure, humidity, wind_speed] for col in numeric_cols: if col in df.columns: df[col] pd.to_numeric(df[col], errorscoerce) # 转换失败则设为NaN return df3.3 数据库操作模块使用SQLAlchemy实现通用存储为了不让代码绑定在某一种数据库上我们使用SQLAlchemy。它通过一个连接字符串Connection String来抽象数据库连接我们只需要更换这个字符串就能轻松切换数据库引擎。首先你需要定义数据表模型。我们可以用SQLAlchemy的Declarative Base来创建一个Python类这个类的属性对应数据库表的列。这里我们定义一个WeatherRecord类。使用declarative_base()创建一个基类然后让我们的模型类继承它。在定义列时要选择合适的数据类型。Integer、Float、String对应数字和字符串DateTime对应时间。primary_keyTrue表示主键。indexTrue可以为经常用于查询的列如station_id和obs_time创建索引这会极大提升查询速度。数据入库时我们使用SQLAlchemy的会话Session来管理。核心方法是session.merge()。merge方法非常智能它会先根据主键这里我们假设是station_id和obs_time的组合在数据库里查找是否存在相同的记录。如果存在就更新那条记录如果不存在就插入一条新记录。这完美解决了数据去重的问题避免了同一时段同一站点的数据被重复插入。实操心得对于大批量数据插入直接使用session.add_all()然后commit效率更高但前提是你确信数据没有重复或者你愿意在插入前先清空旧数据。merge更安全但速度会慢一些因为它需要先查询。对于历史数据回溯填充merge是首选对于实时流式数据追加如果确信不会重复可以用add_all。from sqlalchemy import create_engine, Column, Integer, String, Float, DateTime, Index from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker from datetime import datetime Base declarative_base() class WeatherRecord(Base): 气象数据表模型 __tablename__ weather_data id Column(Integer, primary_keyTrue, autoincrementTrue) # 自增主键方便管理 station_id Column(String(20), nullableFalse) obs_time Column(DateTime, nullableFalse) temperature Column(Float) pressure Column(Float) humidity Column(Float) wind_speed Column(Float) created_at Column(DateTime, defaultdatetime.utcnow) # 创建复合唯一索引确保同一站点同一时间的数据唯一 __table_args__ ( Index(idx_station_time, station_id, obs_time, uniqueTrue), ) class DatabaseManager: def __init__(self, connection_string: str): 初始化数据库连接 connection_string 示例 MySQL: mysqlpymysql://username:passwordlocalhost:3306/weather_db SQLite: sqlite:///./weather_data.db PostgreSQL: postgresqlpsycopg2://username:passwordlocalhost:5432/weather_db self.engine create_engine(connection_string, echoFalse) # echoTrue 会打印所有SQL调试用 # 创建表如果不存在 Base.metadata.create_all(self.engine) Session sessionmaker(bindself.engine) self.Session Session def save_dataframe(self, df: pd.DataFrame, batch_size: int 1000): 将DataFrame数据保存到数据库使用merge实现去重插入 if df.empty: print(DataFrame is empty, nothing to save.) return session self.Session() records_saved 0 try: # 将DataFrame的每一行转换为WeatherRecord对象 for i, row in df.iterrows(): # 处理可能的NaN值SQLAlchemy不接受NaN需要转为None record_dict row.where(pd.notnull(row), None).to_dict() # 构建ORM对象 record WeatherRecord(**record_dict) # 使用merge实现“存在则更新不存在则插入” session.merge(record) # 分批提交避免一次事务过大 if i 0 and i % batch_size 0: session.commit() print(fCommitted {i1} records...) records_saved batch_size session.commit() # 提交剩余记录 total_saved len(df) print(fSuccessfully saved/merged {total_saved} records to database.) except Exception as e: session.rollback() # 发生错误时回滚 print(fError saving data: {e}) raise finally: session.close()4. 完整工作流整合与实战示例现在我们把三个模块像拼积木一样组合起来形成一个完整的、可配置的工作流脚本。这个脚本的入口点会读取配置文件按顺序执行获取数据 - 处理数据 - 存储数据。4.1 配置文件与主程序逻辑我习惯将API密钥、数据库连接等敏感和可配置的信息放在一个单独的配置文件如config.yaml或config.ini里或者从环境变量中读取。主程序main.py负责读取配置、初始化各个模块并控制整个流程。下面是一个使用config.ini和完整主程序的例子config.ini[API] base_url http://api.data.cma.cn/api api_key YOUR_ACTUAL_API_KEY_HERE [DATABASE] # 选择一种数据库连接方式注释掉其他 # sqlite connection_string sqlite:///./weather_data.db # mysql ; connection_string mysqlpymysql://root:passwordlocalhost:3306/weather_db # postgresql ; connection_string postgresqlpsycopg2://postgres:passwordlocalhost:5432/weather_db [QUERY] data_code SURF_CHN_MUL_HOR elements StationID,DATETIME,TEM,PRE,RHU,WIN station_ids 54511, 58362 # 北京上海站示例 start_time 202405010000 end_time 202405012300main.pyimport configparser from pathlib import Path from weather_fetcher import WeatherDataFetcher from data_processor import DataProcessor from database_manager import DatabaseManager, WeatherRecord import pandas as pd def main(): # 1. 读取配置 config configparser.ConfigParser() config.read(config.ini) api_base_url config[API][base_url] api_key config[API][api_key] db_connection_string config[DATABASE][connection_string] # 2. 构建API请求参数 query_params { dataCode: config[QUERY][data_code], elements: config[QUERY][elements], stationIDs: config[QUERY][station_ids], timeRange: f[{config[QUERY][start_time]},{config[QUERY][end_time]}], # 其他可能参数如 orderBy, limit等 orderBy: DATETIME:ASC, # 按时间升序排列 } # 3. 初始化模块 print(Initializing components...) fetcher WeatherDataFetcher(api_keyapi_key, base_urlapi_base_url) db_manager DatabaseManager(connection_stringdb_connection_string) # 4. 执行工作流 print(Fetching data from API...) raw_data fetcher.fetch_data(paramsquery_params) if raw_data: print(Processing data...) df DataProcessor.parse_to_dataframe(raw_data) if not df.empty: print(fData processed, {len(df)} records found.) # 打印前几行看看 print(df.head()) print(Saving to database...) db_manager.save_dataframe(df, batch_size500) print(All done!) else: print(Processed DataFrame is empty. Check the raw data format.) else: print(Failed to fetch data from API.) if __name__ __main__: main()4.2 进阶技巧分页获取与定时任务气象数据网API对单次请求返回的数据量通常有限制比如最多返回1000条。如果你需要获取长时间范围、多站点的数据就需要处理分页。分页逻辑通常依赖于API返回的totalCount总记录数和pageSize每页大小、pageNum页码参数。你需要计算总页数然后用一个循环依次请求每一页的数据并将所有页的数据合并。def fetch_data_with_pagination(fetcher, base_params, page_size1000): 分页获取所有数据 all_records [] page_num 1 base_params[pageSize] page_size while True: base_params[pageNum] page_num print(fFetching page {page_num}...) page_data fetcher.fetch_data(base_params) if not page_data: break records page_data.get(DS, []) if not records: break all_records.extend(records) # 判断是否还有下一页如果返回的记录数小于pageSize通常是最后一页 # 或者根据API返回的totalPage或curPage字段判断 if len(records) page_size: break page_num 1 time.sleep(1) # 礼貌性延迟避免对API服务器造成压力 return {DS: all_records} # 包装成与单次请求相同的格式对于需要定期更新数据的场景比如每天凌晨获取前一天的数据你可以结合操作系统的定时任务如Linux的cron或Windows的“任务计划程序”来运行你的Python脚本。只需将脚本路径和Python解释器路径配置到定时任务中即可。5. 常见问题排查与性能优化在实际运行中你肯定会遇到各种各样的问题。这里我总结几个最常见的问题和解决方法希望能帮你快速排雷。5.1 API请求失败与数据解析错误错误码400或401这通常是请求参数错误或API密钥无效。请仔细检查API密钥是否正确是否已过期。请求的URL和参数名是否与官方文档完全一致注意大小写。时间格式、站点ID格式是否符合要求。你的IP地址是否在API允许的调用白名单内如果API有此限制。错误码429请求过于频繁触发了限流。气象数据网的API通常有调用频率限制如每分钟/每小时最多多少次。解决方案在代码中增加请求间隔使用time.sleep()。如果是批量获取历史数据尽量在夜间等非高峰时段进行。检查是否需要申请更高的调用配额。返回数据为空或DS字段为null首先确认你请求的时间和站点组合确实有数据。有些站点或要素可能在某些时段缺测。检查elements参数是否正确要素代码错误可能导致查询不到数据。使用简单的参数如单个站点、短时间先测试确保基础流程通畅。JSON解析错误偶尔API可能返回非JSON内容如HTML错误页面。在response.json()前可以打印一下response.text的前几百个字符看看或者用try-except捕获JSONDecodeError。5.2 数据库连接与写入问题数据库连接失败检查连接字符串这是最常见的问题。确保用户名、密码、主机地址、端口和数据库名都正确。对于MySQL/PostgreSQL确认数据库服务已启动。检查网络和防火墙确保你的机器可以访问数据库服务器并且对应端口如MySQL的3306是开放的。驱动问题确认已安装正确的数据库驱动pymysql,psycopg2。写入速度慢批量提交如前面代码所示不要每插入一条数据就commit一次而是积累一定数量如1000条后再提交可以大幅减少事务开销。禁用索引仅限大批量初始导入如果你是在初始化一个空表并要导入海量数据如上百万条可以在导入前暂时删除除主键外的索引导入完成后再重建索引。这能极大提升导入速度。但要注意这会使得导入期间表无法用于查询。使用COPY或LOAD DATA对于极大规模数据如数千万条SQLAlchemy的ORM方式可能仍不够快。这时可以考虑直接使用数据库原生的批量导入命令。例如Pandas的DataFrame可以直接用to_sql方法并设置methodmulti或使用if_existsappend但需要注意to_sql在遇到重复主键时可能会报错不如merge灵活。主键或唯一约束冲突这是我们使用session.merge()要解决的核心问题。如果不用merge而直接用add_all就会遇到这个错误。错误信息会明确告诉你违反了哪个唯一约束。检查你的数据中station_id和obs_time的组合是否有重复。5.3 脚本健壮性与日志记录一个用于生产环境的脚本必须有良好的日志记录和异常处理不能只靠print。使用logging模块替换代码中的print语句为不同级别的日志logging.debug,logging.info,logging.warning,logging.error。你可以配置日志输出到文件并设置日志轮转方便日后排查问题。更细致的异常捕获将可能出错的环节网络请求、数据处理、数据库操作用try-except块包裹并记录详细的错误信息包括出错时的参数、数据片段等而不是简单地吞掉异常。添加运行状态检查可以在脚本开始时检查配置文件是否存在、必要参数是否齐全、网络是否可达、数据库是否可连接等提前发现问题。最后我将上述所有模块的代码整合在一个包里并提供了详细的README.md说明文件。你可以在我的代码仓库中找到这个完整的项目。拿到源码后你只需要修改config.ini文件中的API密钥、数据库连接信息和查询参数然后运行python main.py就可以开始自动化获取并存储气象数据了。这套工具我已经稳定运行了半年多处理了数十万条数据希望它也能成为你数据工具箱里的一件利器。如果在使用过程中遇到任何问题或者有更好的改进想法欢迎一起交流。