Loading... # MyBatis-Plus 将“百度天气”写入 PostgreSQL 实践(Spring Boot 3.x) > 目标:以 **Spring Boot 3.x + MyBatis-Plus + WebClient** 拉取“百度天气”JSON并写入 **PostgreSQL**,实现 幂等、可追溯、可扩展 的落库流程。🌦️ ## 一、整体工作流(一图看懂) ```mermaid flowchart LR A[定时任务/手动触发]-->B[WebClient 拉取天气JSON] B-->C[Schema校验&字段映射] C-->D{是否存在相同(city_code,obs_time)?} D--是-->E[ON CONFLICT 更新(幂等)] D--否-->F[插入新记录] E-->G[(PostgreSQL)] F-->G[(PostgreSQL)] G-->H[指标与日志审计 ✅] ``` ## 二、数据库建表(PostgreSQL) ```sql CREATE TABLE IF NOT EXISTS weather_record( id BIGSERIAL PRIMARY KEY, city_code TEXT NOT NULL, city_name TEXT NOT NULL, obs_time TIMESTAMPTZ NOT NULL, temp_c NUMERIC(5,2), weather TEXT, wind_dir TEXT, wind_speed NUMERIC(6,2), humidity INT, pm25 INT, source VARCHAR(32) NOT NULL DEFAULT 'baidu', created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ); CREATE UNIQUE INDEX IF NOT EXISTS unq_weather_city_time ON weather_record(city_code, obs_time); ``` **解释:** * `TIMESTAMPTZ` 统一存储时区;`UNIQUE(city_code, obs_time)` 保证 幂等写入;`source` 标记来源便于稽核。 ## 三、依赖与配置(Spring Boot 3.x) ### 1)`pom.xml` ```xml <dependencies> <!-- Web + JSON --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <!-- MyBatis-Plus --> <dependency> <groupId>com.baomidou</groupId> <artifactId>mybatis-plus-boot-starter</artifactId> <version>3.5.x</version> </dependency> <!-- PostgreSQL 驱动 --> <dependency> <groupId>org.postgresql</groupId> <artifactId>postgresql</artifactId> </dependency> </dependencies> ``` **解释:** * 使用 `mybatis-plus-boot-starter` 集成 MyBatis-Plus;版本用 `3.5.x` 稳定分支;PG 驱动提供 JDBC 连接。 ### 2)`application.yml` ```yaml spring: datasource: url: jdbc:postgresql://127.0.0.1:5432/weatherdb username: weather password: secret hikari: maximum-pool-size: 10 jackson: time-zone: UTC main: banner-mode: off mybatis-plus: configuration: map-underscore-to-camel-case: true global-config: db-config: id-type: AUTO ``` **解释:** * 设定 `UTC` 避免跨区偏移;`map-underscore-to-camel-case` 开启下划线到驼峰映射;`id-type:AUTO` 对应 PG 的 `BIGSERIAL`。 ## 四、实体与 Mapper ### 1)实体 `WeatherRecord.java` ```java @Data @TableName("weather_record") public class WeatherRecord { private Long id; private String cityCode; private String cityName; private OffsetDateTime obsTime; private BigDecimal tempC; private String weather; private String windDir; private BigDecimal windSpeed; private Integer humidity; private Integer pm25; private String source; private OffsetDateTime createdAt; private OffsetDateTime updatedAt; } ``` **解释:** * 字段与表结构一一对应;`OffsetDateTime` 与 `TIMESTAMPTZ` 匹配,时区安全。 ### 2)Mapper(带 UPSERT)`WeatherRecordMapper.java` ```java @Mapper public interface WeatherRecordMapper extends BaseMapper<WeatherRecord> { @Insert(""" INSERT INTO weather_record (city_code,city_name,obs_time,temp_c,weather,wind_dir,wind_speed,humidity,pm25,source,created_at,updated_at) VALUES (#{cityCode},#{cityName},#{obsTime},#{tempC},#{weather},#{windDir},#{windSpeed},#{humidity},#{pm25},#{source},NOW(),NOW()) ON CONFLICT (city_code, obs_time) DO UPDATE SET city_name = EXCLUDED.city_name, temp_c = EXCLUDED.temp_c, weather = EXCLUDED.weather, wind_dir = EXCLUDED.wind_dir, wind_speed= EXCLUDED.wind_speed, humidity = EXCLUDED.humidity, pm25 = EXCLUDED.pm25, source = EXCLUDED.source, updated_at= NOW() """) int upsert(WeatherRecord r); } ``` **解释:** * 使用 PG 的 `ON CONFLICT` 实现 幂等写入;只在唯一键冲突时更新最新数据;返回影响行数便于统计。 ## 五、拉取与落库(Service) ```java @Service @RequiredArgsConstructor public class WeatherIngestService { private final WebClient webClient = WebClient.builder().build(); private final WeatherRecordMapper mapper; @Value("${baidu.weather.endpoint}") String endpoint; // 例如官方天气REST接口 @Value("${baidu.ak}") String ak; public int ingest(String cityCode, String cityName) { // 1) 拉取 JSON(参数示意:城市编码/名称、AK 等) Map<String, Object> json = webClient.get() .uri(uri -> uri.path(endpoint) .queryParam("city", cityName) .queryParam("code", cityCode) .queryParam("ak", ak).build()) .retrieve().bodyToMono(new ParameterizedTypeReference<Map<String,Object>>(){}).block(); // 2) 解析映射(按实际JSON结构调整键名) Map<String, Object> now = (Map<String, Object>) ((Map<String,Object>)json.get("result")).get("now"); WeatherRecord r = new WeatherRecord(); r.setCityCode(cityCode); r.setCityName(cityName); r.setObsTime(OffsetDateTime.parse(now.get("obsTime").toString())); // ISO8601 r.setTempC(new BigDecimal(now.get("temp").toString())); r.setWeather(now.get("text").toString()); r.setWindDir(now.get("windDir").toString()); r.setWindSpeed(new BigDecimal(now.get("windSpeed").toString())); r.setHumidity(Integer.valueOf(now.get("rh").toString())); r.setPm25(now.get("pm25")==null?null:Integer.valueOf(now.get("pm25").toString())); r.setSource("baidu"); // 3) 落库(UPSERT) return mapper.upsert(r); } } ``` **解释:** * `WebClient` 发起 GET;参数名与 JSON 键需按照你的实际接口定义替换; * `OffsetDateTime.parse` 解析 ISO8601 时间戳,确保 时间精度与时区一致; * 最终调用 `upsert` 保证 **重复采集不产生脏数据**。 ## 六、触发入口与定时任务 ```java @RestController @RequiredArgsConstructor @RequestMapping("/ingest") public class IngestController { private final WeatherIngestService svc; @PostMapping("/now") public String now(@RequestParam String cityCode, @RequestParam String cityName){ int n = svc.ingest(cityCode, cityName); return "affected=" + n; } } @Component @RequiredArgsConstructor @EnableScheduling class WeatherScheduler { private final WeatherIngestService svc; // 每 10 分钟采集一次(示例) @Scheduled(cron = "0 */10 * * * *") public void cron(){ svc.ingest("110000","北京"); svc.ingest("310000","上海"); } } ``` **解释:** * 提供 `/ingest/now` 便于手动触发; * `@Scheduled` 定时采集关键城市,实现 自动化。 ## 七、字段映射与质量控制(对照表) | 采集字段 | 示例 | 目标列 | 备注 | | ----------- | ----------------------------- | --------------------------- | -------------------- | | city / code | “北京”/“110000” | `city_name`/`city_code` | 城市名与行政编码 | | obsTime | `2025-08-29T10:30:00+08:00` | `obs_time` | 存为 `TIMESTAMPTZ` | | temp | `33.2` | `temp_c` | 摄氏温度 | | text | “多云” | `weather` | 天气现象 | | windDir | “东北风” | `wind_dir` | 风向 | | windSpeed | `3.5` | `wind_speed` | m/s 或 km/h,需统一 | | rh | `60` | `humidity` | 相对湿度 | | pm25 | `35` | `pm25` | 可能缺失,注意判空 | **解释:** 通过表格把源字段与目标列 一一对应,避免上线后口径不一致。 ## 八、幂等与告警的“可解释公式” ```text 幂等键 = (city_code, obs_time) 若 同一键重复到达 → ON CONFLICT 更新最新值 ``` ```text 告警条件: 1) 连续3个周期 affected = 0 → 可能接口异常 2) 连续N条 obs_time 落后当前时间 > Δt → 可能卡住 ``` **解释:** 用 唯一键 与 时序阈值 做健康度判断,简洁可靠。 ## 九、常见坑与加固 ✅ * 时区统一:接口时间若为本地时,需转为 UTC 后入库。 * 单位统一:风速/温度单位必须明确(m/s vs km/h)。 * 缺失字段:如 PM2.5 为空,优先存 `NULL` 而非 0,避免误解读。 * 重试与熔断:网络失败应有限次重试;长期失败应熔断+告警。 * 速率限制:遵守服务侧 QPS 要求,防止被限流。 --- ### 小结 通过 **MyBatis-Plus UPSERT**、**PostgreSQL 唯一键** 与 **WebClient 拉取**,即可快速完成“百度天气”到 PG 的 稳定、幂等 写入。表结构清晰、字段映射明确、时间与单位统一,是保障数据质量的关键。🚀 最后修改:2025 年 09 月 30 日 © 允许规范转载 打赏 赞赏作者 支付宝微信 赞 如果觉得我的文章对你有用,请随意赞赏