diff --git a/app/api/v1/endpoints/timeseries/composite.py b/app/api/v1/endpoints/timeseries/composite.py index c0cbf41..b863097 100644 --- a/app/api/v1/endpoints/timeseries/composite.py +++ b/app/api/v1/endpoints/timeseries/composite.py @@ -8,8 +8,7 @@ from .dependencies import get_timescale_connection, get_postgres_connection router = APIRouter() -@router.get("/composite/scada-simulation", summary="获取SCADA关联的模拟数据", - tags=["复合查询"]) +@router.get("/composite/scada-simulation", summary="获取SCADA关联的模拟数据") async def get_scada_associated_simulation_data( start_time: datetime = Query(..., description="查询开始时间"), end_time: datetime = Query(..., description="查询结束时间"), @@ -74,8 +73,7 @@ async def get_scada_associated_simulation_data( raise HTTPException(status_code=400, detail=str(e)) -@router.get("/composite/element-simulation", summary="获取管网元素的模拟数据", - tags=["复合查询"]) +@router.get("/composite/element-simulation", summary="获取管网元素的模拟数据") async def get_feature_simulation_data( start_time: datetime = Query(..., description="查询开始时间"), end_time: datetime = Query(..., description="查询结束时间"), @@ -145,8 +143,7 @@ async def get_feature_simulation_data( raise HTTPException(status_code=400, detail=str(e)) -@router.get("/composite/element-scada", summary="获取管网元素关联的SCADA监测数据", - tags=["复合查询"]) +@router.get("/composite/element-scada", summary="获取管网元素关联的SCADA监测数据") async def get_element_associated_scada_data( element_id: str = Query(..., description="管网元素ID(管道或节点)"), start_time: datetime = Query(..., description="查询开始时间"), @@ -188,8 +185,7 @@ async def get_element_associated_scada_data( raise HTTPException(status_code=400, detail=str(e)) -@router.post("/composite/clean-scada", summary="清洗SCADA监测数据", - tags=["复合查询"]) +@router.post("/composite/clean-scada", summary="清洗SCADA监测数据") async def clean_scada_data( device_ids: str = Query(..., description="设备ID列表或 'all' 表示清洗所有设备"), start_time: datetime = Query(..., description="清洗数据的开始时间"), @@ -232,8 +228,7 @@ async def clean_scada_data( raise HTTPException(status_code=400, detail=str(e)) -@router.get("/composite/pipeline-health-prediction", summary="预测管道健康状况", - tags=["复合查询"]) +@router.get("/composite/pipeline-health-prediction", summary="预测管道健康状况") async def predict_pipeline_health( query_time: datetime = Query(..., description="查询时间"), network_name: str = Query(..., description="管网名称(或数据库名称)"), diff --git a/app/api/v1/endpoints/timeseries/realtime.py b/app/api/v1/endpoints/timeseries/realtime.py index d6fabf3..eb6ee8b 100644 --- a/app/api/v1/endpoints/timeseries/realtime.py +++ b/app/api/v1/endpoints/timeseries/realtime.py @@ -8,9 +8,12 @@ from .dependencies import get_timescale_connection router = APIRouter() +TIME_WITH_TZ_DESC = "ISO 8601 / RFC 3339 时间,必须显式带时区;可直接传 UTC+8,服务端会先转换为 UTC 再处理。" +TIME_RANGE_START_DESC = f"时间范围开始时间。{TIME_WITH_TZ_DESC}" +TIME_RANGE_END_DESC = f"时间范围结束时间。{TIME_WITH_TZ_DESC}" -@router.post("/realtime/links/batch", status_code=201, summary="批量插入实时管道数据", - tags=["时间序列-实时数据"]) + +@router.post("/realtime/links/batch", status_code=201, summary="批量插入实时管道数据") async def insert_realtime_links( data: List[dict] = Body(..., description="管道数据列表,每项包含管道ID、时间戳等信息"), conn: AsyncConnection = Depends(get_timescale_connection) @@ -30,16 +33,21 @@ async def insert_realtime_links( return {"message": f"Inserted {len(data)} records"} -@router.get("/realtime/links", summary="查询实时管道数据", tags=["时间序列-实时数据"]) +@router.get( + "/realtime/links", + summary="查询实时管道数据", + description="按时间范围查询实时管道数据。start_time 和 end_time 必须显式带时区;允许传 UTC+8,服务端会先归一化为 UTC 再执行查询。", +) async def get_realtime_links( - start_time: datetime = Query(..., description="查询开始时间"), - end_time: datetime = Query(..., description="查询结束时间"), + start_time: datetime = Query(..., description=TIME_RANGE_START_DESC), + end_time: datetime = Query(..., description=TIME_RANGE_END_DESC), conn: AsyncConnection = Depends(get_timescale_connection), ): """ 查询指定时间范围内的实时管道数据 - 根据时间范围查询所有实时管道的监测值。 + 根据时间范围查询所有实时管道的监测值。传入时间必须显式包含时区, + 可以直接使用 UTC+8,服务端会先统一转换为 UTC 再参与数据库查询。 Args: start_time: 查询开始时间 @@ -51,10 +59,14 @@ async def get_realtime_links( return await RealtimeRepository.get_links_by_time_range(conn, start_time, end_time) -@router.delete("/realtime/links", summary="删除实时管道数据", tags=["时间序列-实时数据"]) +@router.delete( + "/realtime/links", + summary="删除实时管道数据", + description="按时间范围删除实时管道数据。start_time 和 end_time 必须显式带时区;允许传 UTC+8,服务端按请求中的绝对时间删除对应 UTC 数据。", +) async def delete_realtime_links( - start_time: datetime = Query(..., description="删除开始时间"), - end_time: datetime = Query(..., description="删除结束时间"), + start_time: datetime = Query(..., description=TIME_RANGE_START_DESC), + end_time: datetime = Query(..., description=TIME_RANGE_END_DESC), conn: AsyncConnection = Depends(get_timescale_connection), ): """ @@ -73,11 +85,10 @@ async def delete_realtime_links( return {"message": "Deleted successfully"} -@router.patch("/realtime/links/{link_id}/field", summary="更新实时管道字段", - tags=["时间序列-实时数据"]) +@router.patch("/realtime/links/{link_id}/field", summary="更新实时管道字段") async def update_realtime_link_field( link_id: str = Path(..., description="管道ID"), - time: datetime = Query(..., description="更新数据的时间戳"), + time: datetime = Query(..., description=f"要更新记录的时间戳。{TIME_WITH_TZ_DESC}"), field: str = Query(..., description="要更新的字段名称"), value: float = Query(..., description="更新的字段值"), conn: AsyncConnection = Depends(get_timescale_connection), @@ -106,8 +117,7 @@ async def update_realtime_link_field( raise HTTPException(status_code=400, detail=str(e)) -@router.post("/realtime/nodes/batch", status_code=201, summary="批量插入实时节点数据", - tags=["时间序列-实时数据"]) +@router.post("/realtime/nodes/batch", status_code=201, summary="批量插入实时节点数据") async def insert_realtime_nodes( data: List[dict] = Body(..., description="节点数据列表,每项包含节点ID、时间戳等信息"), conn: AsyncConnection = Depends(get_timescale_connection) @@ -127,16 +137,21 @@ async def insert_realtime_nodes( return {"message": f"Inserted {len(data)} records"} -@router.get("/realtime/nodes", summary="查询实时节点数据", tags=["时间序列-实时数据"]) +@router.get( + "/realtime/nodes", + summary="查询实时节点数据", + description="按时间范围查询实时节点数据。start_time 和 end_time 必须显式带时区;允许传 UTC+8,服务端会先归一化为 UTC 再执行查询。", +) async def get_realtime_nodes( - start_time: datetime = Query(..., description="查询开始时间"), - end_time: datetime = Query(..., description="查询结束时间"), + start_time: datetime = Query(..., description=TIME_RANGE_START_DESC), + end_time: datetime = Query(..., description=TIME_RANGE_END_DESC), conn: AsyncConnection = Depends(get_timescale_connection), ): """ 查询指定时间范围内的实时节点数据 - 根据时间范围查询所有实时节点的监测值。 + 根据时间范围查询所有实时节点的监测值。传入时间必须显式包含时区, + 可以直接使用 UTC+8,服务端会先统一转换为 UTC 再参与数据库查询。 Args: start_time: 查询开始时间 @@ -148,10 +163,14 @@ async def get_realtime_nodes( return await RealtimeRepository.get_nodes_by_time_range(conn, start_time, end_time) -@router.delete("/realtime/nodes", summary="删除实时节点数据", tags=["时间序列-实时数据"]) +@router.delete( + "/realtime/nodes", + summary="删除实时节点数据", + description="按时间范围删除实时节点数据。start_time 和 end_time 必须显式带时区;允许传 UTC+8,服务端按请求中的绝对时间删除对应 UTC 数据。", +) async def delete_realtime_nodes( - start_time: datetime = Query(..., description="删除开始时间"), - end_time: datetime = Query(..., description="删除结束时间"), + start_time: datetime = Query(..., description=TIME_RANGE_START_DESC), + end_time: datetime = Query(..., description=TIME_RANGE_END_DESC), conn: AsyncConnection = Depends(get_timescale_connection), ): """ @@ -172,12 +191,11 @@ async def delete_realtime_nodes( -@router.post("/realtime/simulation/store", status_code=201, summary="存储实时模拟结果", - tags=["时间序列-实时数据"]) +@router.post("/realtime/simulation/store", status_code=201, summary="存储实时模拟结果") async def store_realtime_simulation_result( node_result_list: List[dict] = Body(..., description="节点模拟结果列表"), link_result_list: List[dict] = Body(..., description="管道模拟结果列表"), - result_start_time: str = Query(..., description="模拟结果开始时间"), + result_start_time: str = Query(..., description=f"模拟结果开始时间。{TIME_WITH_TZ_DESC}"), conn: AsyncConnection = Depends(get_timescale_connection), ): """ @@ -199,10 +217,13 @@ async def store_realtime_simulation_result( return {"message": "Simulation results stored successfully"} -@router.get("/realtime/query/by-time-property", summary="按时间和属性查询实时数据", - tags=["时间序列-实时数据"]) +@router.get( + "/realtime/query/by-time-property", + summary="按时间和属性查询实时数据", + description="查询指定时间点的实时属性值。query_time 必须显式带时区;允许传 UTC+8,服务端会先归一化为 UTC 再执行查询。", +) async def query_realtime_records_by_time_property( - query_time: str = Query(..., description="查询时间"), + query_time: str = Query(..., description=f"查询时间。{TIME_WITH_TZ_DESC}"), type: str = Query(..., description="数据类型,pipe(管道)或 junction(节点)"), property: str = Query(..., description="要查询的属性名称"), conn: AsyncConnection = Depends(get_timescale_connection), @@ -232,12 +253,15 @@ async def query_realtime_records_by_time_property( raise HTTPException(status_code=400, detail=str(e)) -@router.get("/realtime/query/by-id-time", summary="按ID和时间查询实时模拟数据", - tags=["时间序列-实时数据"]) +@router.get( + "/realtime/query/by-id-time", + summary="按ID和时间查询实时模拟数据", + description="查询指定元素在某一时间点的实时模拟结果。query_time 必须显式带时区;允许传 UTC+8,服务端会先归一化为 UTC 再执行查询。", +) async def query_realtime_simulation_by_id_time( id: str = Query(..., description="元素ID(管道ID或节点ID)"), type: str = Query(..., description="元素类型,pipe(管道)或 junction(节点)"), - query_time: str = Query(..., description="查询时间"), + query_time: str = Query(..., description=f"查询时间。{TIME_WITH_TZ_DESC}"), conn: AsyncConnection = Depends(get_timescale_connection), ): """ diff --git a/app/api/v1/endpoints/timeseries/scada.py b/app/api/v1/endpoints/timeseries/scada.py index a42f87c..3ed1bab 100644 --- a/app/api/v1/endpoints/timeseries/scada.py +++ b/app/api/v1/endpoints/timeseries/scada.py @@ -9,20 +9,19 @@ from .dependencies import get_timescale_connection router = APIRouter() -@router.post("/scada/batch", status_code=201, summary="批量插入SCADA监测数据", - tags=["时间序列-监测数据"]) +@router.post("/scada/batch", status_code=201, summary="批量插入SCADA监测数据") async def insert_scada_data( data: List[dict] = Body(..., description="SCADA设备监测数据列表"), - conn: AsyncConnection = Depends(get_timescale_connection) + conn: AsyncConnection = Depends(get_timescale_connection), ): """ 批量插入SCADA监测数据 - + 将多个设备的实时监测数据批量插入时间序列数据库。 - + Args: data: SCADA设备监测数据列表,每项包含device_id、时间戳和监测值等信息 - + Returns: 插入成功的记录数 """ @@ -30,24 +29,25 @@ async def insert_scada_data( return {"message": f"Inserted {len(data)} records"} -@router.get("/scada/by-ids-time-range", summary="按设备ID和时间范围查询SCADA数据", - tags=["时间序列-监测数据"]) +@router.get("/scada/by-ids-time-range", summary="按设备ID和时间范围查询SCADA数据") async def get_scada_by_ids_time_range( start_time: datetime = Query(..., description="查询开始时间"), end_time: datetime = Query(..., description="查询结束时间"), - device_ids: str = Query(..., description="设备ID列���,逗号分隔,如 'device1,device2,device3'"), + device_ids: str = Query( + ..., description="设备ID列表,逗号分隔,如 'device1,device2,device3'" + ), conn: AsyncConnection = Depends(get_timescale_connection), ): """ 按设备ID和时间范围查询SCADA监测数据 - + 查询多个设备在指定时间范围内的所有监测数据。 - + Args: start_time: 查询开始时间 end_time: 查询结束时间 device_ids: 设备ID列表,用逗号分隔 - + Returns: SCADA监测数据列表 """ @@ -59,29 +59,32 @@ async def get_scada_by_ids_time_range( ) -@router.get("/scada/by-ids-field-time-range", summary="按设备ID、字段和时间范围查询SCADA数据", - tags=["时间序列-监测数据"]) +@router.get( + "/scada/by-ids-field-time-range", summary="按设备ID、字段和时间范围查询SCADA数据" +) async def get_scada_field_by_ids_time_range( start_time: datetime = Query(..., description="查询开始时间"), end_time: datetime = Query(..., description="查询结束时间"), field: str = Query(..., description="要查询的字段名称"), - device_ids: str = Query(..., description="设备ID列表,逗号分隔,如 'device1,device2,device3'"), + device_ids: str = Query( + ..., description="设备ID列表,逗号分隔,如 'device1,device2,device3'" + ), conn: AsyncConnection = Depends(get_timescale_connection), ): """ 按设备ID、字段和时间范围查询特定SCADA数据 - + 查询多个设备在指定时间范围内的特定字段监测数据。 - + Args: start_time: 查询开始时间 end_time: 查询结束时间 field: 字段名称 device_ids: 设备ID列表,用逗号分隔 - + Returns: SCADA字段数据列表 - + Raises: HTTPException: 当字段不存在或查询参数无效时返回400错误 """ @@ -98,8 +101,7 @@ async def get_scada_field_by_ids_time_range( raise HTTPException(status_code=400, detail=str(e)) -@router.patch("/scada/{device_id}/field", summary="更新SCADA设备字段", - tags=["时间序列-监测数据"]) +@router.patch("/scada/{device_id}/field", summary="更新SCADA设备字段") async def update_scada_field( device_id: str = Path(..., description="设备ID"), time: datetime = Query(..., description="更新数据的时间戳"), @@ -109,18 +111,18 @@ async def update_scada_field( ): """ 更新指定设备的字段值 - + 更新SCADA设备在特定时间的某个字段监测数据。 - + Args: device_id: 设备ID time: 数据时间戳 field: 字段名称 value: 字段新值 - + Returns: 更新结果信息 - + Raises: HTTPException: 当字段不存在或更新失败时返回400错误 """ @@ -131,8 +133,7 @@ async def update_scada_field( raise HTTPException(status_code=400, detail=str(e)) -@router.delete("/scada/by-id-time-range", summary="按设备ID和时间范围删除SCADA数据", - tags=["时间序列-监测数据"]) +@router.delete("/scada/by-id-time-range", summary="按设备ID和时间范围删除SCADA数据") async def delete_scada_data( device_id: str = Query(..., description="设备ID"), start_time: datetime = Query(..., description="删除开始时间"), @@ -141,14 +142,14 @@ async def delete_scada_data( ): """ 删除指定设备和时间范围内的SCADA数据 - + 删除在指定时间范围内的特定设备监测数据。 - + Args: device_id: 设备ID start_time: 删除开始时间 end_time: 删除结束时间 - + Returns: 删除结果信息 """ diff --git a/app/api/v1/endpoints/timeseries/scheme.py b/app/api/v1/endpoints/timeseries/scheme.py index 7e342c0..76f71e6 100644 --- a/app/api/v1/endpoints/timeseries/scheme.py +++ b/app/api/v1/endpoints/timeseries/scheme.py @@ -9,20 +9,19 @@ from .dependencies import get_timescale_connection router = APIRouter() -@router.post("/scheme/links/batch", status_code=201, summary="批量插入方案管道数据", - tags=["时间序列-方案数据"]) +@router.post("/scheme/links/batch", status_code=201, summary="批量插入方案管道数据") async def insert_scheme_links( data: List[dict] = Body(..., description="方案管道数据列表"), - conn: AsyncConnection = Depends(get_timescale_connection) + conn: AsyncConnection = Depends(get_timescale_connection), ): """ 批量插入方案管道数据 - + 将特定方案的管道模拟数据批量插入时间序列数据库。 - + Args: data: 方案管道数据列表 - + Returns: 插入成功的记录数 """ @@ -30,7 +29,7 @@ async def insert_scheme_links( return {"message": f"Inserted {len(data)} records"} -@router.get("/scheme/links", summary="查询方案管道数据", tags=["时间序列-方案数据"]) +@router.get("/scheme/links", summary="查询方案管道数据") async def get_scheme_links( scheme_type: str = Query(..., description="方案类型"), scheme_name: str = Query(..., description="方案名称"), @@ -40,15 +39,15 @@ async def get_scheme_links( ): """ 查询指定方案和时间范围内的管道数据 - + 根据方案和时间范围查询管道的模拟值。 - + Args: scheme_type: 方案类型 scheme_name: 方案名称 start_time: 查询开始时间 end_time: 查询结束时间 - + Returns: 方案管道数据列表 """ @@ -57,8 +56,7 @@ async def get_scheme_links( ) -@router.get("/scheme/links/{link_id}/field", summary="查询方案管道字段数据", - tags=["时间序列-方案数据"]) +@router.get("/scheme/links/{link_id}/field", summary="查询方案管道字段数据") async def get_scheme_link_field( link_id: str = Path(..., description="管道ID"), scheme_type: str = Query(..., description="方案类型"), @@ -70,9 +68,9 @@ async def get_scheme_link_field( ): """ 查询指定方案管道的特定字段数据 - + 查询特定方案中指定管道在时间范围内的特定字段值。 - + Args: link_id: 管道ID scheme_type: 方案类型 @@ -80,10 +78,10 @@ async def get_scheme_link_field( start_time: 查询开始时间 end_time: 查询结束时间 field: 字段名称 - + Returns: 字段数据列表 - + Raises: HTTPException: 当查询参数无效时返回400错误 """ @@ -95,8 +93,7 @@ async def get_scheme_link_field( raise HTTPException(status_code=400, detail=str(e)) -@router.patch("/scheme/links/{link_id}/field", summary="更新方案管道字段", - tags=["时间序列-方案数据"]) +@router.patch("/scheme/links/{link_id}/field", summary="更新方案管道字段") async def update_scheme_link_field( link_id: str = Path(..., description="管道ID"), scheme_type: str = Query(..., description="方案类型"), @@ -108,9 +105,9 @@ async def update_scheme_link_field( ): """ 更新指定方案管道的字段值 - + 更新特定方案中指定管道在某个时间的字段数据。 - + Args: link_id: 管道ID scheme_type: 方案类型 @@ -118,10 +115,10 @@ async def update_scheme_link_field( time: 数据时间戳 field: 字段名称 value: 字段新值 - + Returns: 更新结果信息 - + Raises: HTTPException: 当字段不存在或更新失败时返回400错误 """ @@ -134,7 +131,7 @@ async def update_scheme_link_field( raise HTTPException(status_code=400, detail=str(e)) -@router.delete("/scheme/links", summary="删除方案管道数据", tags=["时间序列-方案数据"]) +@router.delete("/scheme/links", summary="删除方案管道数据") async def delete_scheme_links( scheme_type: str = Query(..., description="方案类型"), scheme_name: str = Query(..., description="方案名称"), @@ -144,15 +141,15 @@ async def delete_scheme_links( ): """ 删除指定方案和时间范围内的管道数据 - + 删除在指定方案和时间范围内的所有管道模拟数据。 - + Args: scheme_type: 方案类型 scheme_name: 方案名称 start_time: 删除开始时间 end_time: 删除结束时间 - + Returns: 删除结果信息 """ @@ -162,20 +159,19 @@ async def delete_scheme_links( return {"message": "Deleted successfully"} -@router.post("/scheme/nodes/batch", status_code=201, summary="批量插入方案节点数据", - tags=["时间序列-方案数据"]) +@router.post("/scheme/nodes/batch", status_code=201, summary="批量插入方案节点数据") async def insert_scheme_nodes( data: List[dict] = Body(..., description="方案节点数据列表"), - conn: AsyncConnection = Depends(get_timescale_connection) + conn: AsyncConnection = Depends(get_timescale_connection), ): """ 批量插入方案节点数据 - + 将特定方案的节点模拟数据批量插入时间序列数据库。 - + Args: data: 方案节点数据列表 - + Returns: 插入成功的记录数 """ @@ -183,8 +179,7 @@ async def insert_scheme_nodes( return {"message": f"Inserted {len(data)} records"} -@router.get("/scheme/nodes/{node_id}/field", summary="查询方案节点字段数据", - tags=["时间序列-方案数据"]) +@router.get("/scheme/nodes/{node_id}/field", summary="查询方案节点字段数据") async def get_scheme_node_field( node_id: str = Path(..., description="节点ID"), scheme_type: str = Query(..., description="方案类型"), @@ -196,9 +191,9 @@ async def get_scheme_node_field( ): """ 查询指定方案节点的特定字段数据 - + 查询特定方案中指定节点在时间范围内的特定字段值。 - + Args: node_id: 节点ID scheme_type: 方案类型 @@ -206,10 +201,10 @@ async def get_scheme_node_field( start_time: 查询开始时间 end_time: 查询结束时间 field: 字段名称 - + Returns: 字段数据列表 - + Raises: HTTPException: 当查询参数无效时返回400错误 """ @@ -221,8 +216,7 @@ async def get_scheme_node_field( raise HTTPException(status_code=400, detail=str(e)) -@router.patch("/scheme/nodes/{node_id}/field", summary="更新方案节点字段", - tags=["时间序列-方案数据"]) +@router.patch("/scheme/nodes/{node_id}/field", summary="更新方案节点字段") async def update_scheme_node_field( node_id: str = Path(..., description="节点ID"), scheme_type: str = Query(..., description="方案类型"), @@ -234,9 +228,9 @@ async def update_scheme_node_field( ): """ 更新指定方案节点的字段值 - + 更新特定方案中指定节点在某个时间的字段数据。 - + Args: node_id: 节点ID scheme_type: 方案类型 @@ -244,10 +238,10 @@ async def update_scheme_node_field( time: 数据时间戳 field: 字段名称 value: 字段新值 - + Returns: 更新结果信息 - + Raises: HTTPException: 当字段不存在或更新失败时返回400错误 """ @@ -260,7 +254,7 @@ async def update_scheme_node_field( raise HTTPException(status_code=400, detail=str(e)) -@router.delete("/scheme/nodes", summary="删除方案节点数据", tags=["时间序列-方案数据"]) +@router.delete("/scheme/nodes", summary="删除方案节点数据") async def delete_scheme_nodes( scheme_type: str = Query(..., description="方案类型"), scheme_name: str = Query(..., description="方案名称"), @@ -270,15 +264,15 @@ async def delete_scheme_nodes( ): """ 删除指定方案和时间范围内的节点数据 - + 删除在指定方案和时间范围内的所有节点模拟数据。 - + Args: scheme_type: 方案类型 scheme_name: 方案名称 start_time: 删除开始时间 end_time: 删除结束时间 - + Returns: 删除结果信息 """ @@ -288,8 +282,7 @@ async def delete_scheme_nodes( return {"message": "Deleted successfully"} -@router.post("/scheme/simulation/store", status_code=201, summary="存储方案模拟结果", - tags=["时间序列-方案数据"]) +@router.post("/scheme/simulation/store", status_code=201, summary="存储方案模拟结果") async def store_scheme_simulation_result( scheme_type: str = Query(..., description="方案类型"), scheme_name: str = Query(..., description="方案名称"), @@ -300,16 +293,16 @@ async def store_scheme_simulation_result( ): """ 存储方案模拟结果到时间序列数据库 - + 将特定方案的节点和管道模拟计算结果批量存储到TimescaleDB数据库。 - + Args: scheme_type: 方案类型 scheme_name: 方案名称 node_result_list: 节点模拟结果列表 link_result_list: 管道模拟结果列表 result_start_time: 模拟结果对应的起始时间 - + Returns: 存储结果信息 """ @@ -324,8 +317,9 @@ async def store_scheme_simulation_result( return {"message": "Scheme simulation results stored successfully"} -@router.get("/scheme/query/by-scheme-time-property", summary="按方案、时间和属性查询数据", - tags=["时间序列-方案数据"]) +@router.get( + "/scheme/query/by-scheme-time-property", summary="按方案、时间和属性查询数据" +) async def query_scheme_records_by_scheme_time_property( scheme_type: str = Query(..., description="方案类型"), scheme_name: str = Query(..., description="方案名称"), @@ -336,19 +330,19 @@ async def query_scheme_records_by_scheme_time_property( ): """ 按指定方案、时间和属性查询所有方案数据 - + 查询在特定方案和时间点,所有指定类型元素的特定属性值。 - + Args: scheme_type: 方案类型 scheme_name: 方案名称 query_time: 查询时间 type: 元素类型(pipe或junction) property: 属性名称 - + Returns: 查询结果列表 - + Raises: HTTPException: 当查询参数无效时返回400错误 """ @@ -361,8 +355,7 @@ async def query_scheme_records_by_scheme_time_property( raise HTTPException(status_code=400, detail=str(e)) -@router.get("/scheme/query/by-id-time", summary="按ID和时间查询方案模拟数据", - tags=["时间序列-方案数据"]) +@router.get("/scheme/query/by-id-time", summary="按ID和时间查询方案模拟数据") async def query_scheme_simulation_by_id_time( scheme_type: str = Query(..., description="方案类型"), scheme_name: str = Query(..., description="方案名称"), @@ -373,19 +366,19 @@ async def query_scheme_simulation_by_id_time( ): """ 按指定ID和时间查询方案模拟结果 - + 查询特定方案中的元素在某一时间点的模拟数据。 - + Args: scheme_type: 方案类型 scheme_name: 方案名称 id: 元素ID type: 元素类型(pipe或junction) query_time: 查询时间 - + Returns: 模拟结果数据 - + Raises: HTTPException: 当查询参数无效时返回400错误 """ diff --git a/app/infra/db/timescaledb/internal_queries.py b/app/infra/db/timescaledb/internal_queries.py index 3de4db7..fda46f3 100644 --- a/app/infra/db/timescaledb/internal_queries.py +++ b/app/infra/db/timescaledb/internal_queries.py @@ -10,6 +10,7 @@ from app.core.config import get_timescaledb_pgconn_string from app.infra.db.timescaledb.repositories.scheme import SchemeRepository from app.infra.db.timescaledb.repositories.realtime import RealtimeRepository from app.infra.db.timescaledb.repositories.scada import ScadaRepository +from app.services.time_api import parse_utc_time class InternalStorage: @@ -89,10 +90,9 @@ class InternalQueries: ) -> dict: """查询指定时间点的 SCADA 数据""" - # 解析时间,假设是北京时间 - beijing_time = datetime.fromisoformat(query_time) - start_time = beijing_time - timedelta(seconds=1) - end_time = beijing_time + timedelta(seconds=1) + target_time = parse_utc_time(query_time, field_name="query_time") + start_time = target_time - timedelta(seconds=1) + end_time = target_time + timedelta(seconds=1) for attempt in range(max_retries): try: @@ -132,14 +132,8 @@ class InternalQueries: max_retries: int = 3, ) -> dict[str, list[dict]]: """查询指定时间窗的 SCADA 数据,返回 {device_id: [{time, value}, ...]}。""" - start_dt = ( - datetime.fromisoformat(start_time) - if isinstance(start_time, str) - else start_time - ) - end_dt = ( - datetime.fromisoformat(end_time) if isinstance(end_time, str) else end_time - ) + start_dt = parse_utc_time(start_time, field_name="start_time") + end_dt = parse_utc_time(end_time, field_name="end_time") for attempt in range(max_retries): try: @@ -238,14 +232,8 @@ class InternalQueries: if not element_ids: return {} - start_dt = ( - datetime.fromisoformat(start_time) - if isinstance(start_time, str) - else start_time - ) - end_dt = ( - datetime.fromisoformat(end_time) if isinstance(end_time, str) else end_time - ) + start_dt = parse_utc_time(start_time, field_name="start_time") + end_dt = parse_utc_time(end_time, field_name="end_time") table_name, valid_fields = InternalQueries._resolve_simulation_table(element_type) if field not in valid_fields: raise ValueError(f"Invalid field for {element_type}: {field}") diff --git a/app/infra/db/timescaledb/repositories/realtime.py b/app/infra/db/timescaledb/repositories/realtime.py index 06a32de..6b26fa4 100644 --- a/app/infra/db/timescaledb/repositories/realtime.py +++ b/app/infra/db/timescaledb/repositories/realtime.py @@ -1,10 +1,8 @@ from typing import List, Any, Dict -from datetime import datetime, timedelta, timezone +from datetime import datetime, timedelta from collections import defaultdict from psycopg import AsyncConnection, Connection, sql - -# 定义UTC+8时区 -UTC_8 = timezone(timedelta(hours=8)) +from app.services.time_api import parse_utc_time class RealtimeRepository: @@ -102,10 +100,12 @@ class RealtimeRepository: async def get_links_by_time_range( conn: AsyncConnection, start_time: datetime, end_time: datetime ) -> List[dict]: + normalized_start_time = parse_utc_time(start_time, field_name="start_time") + normalized_end_time = parse_utc_time(end_time, field_name="end_time") async with conn.cursor() as cur: await cur.execute( "SELECT * FROM realtime.link_simulation WHERE time >= %s AND time <= %s", - (start_time, end_time), + (normalized_start_time, normalized_end_time), ) return await cur.fetchall() @@ -298,10 +298,12 @@ class RealtimeRepository: async def get_nodes_by_time_range( conn: AsyncConnection, start_time: datetime, end_time: datetime ) -> List[dict]: + normalized_start_time = parse_utc_time(start_time, field_name="start_time") + normalized_end_time = parse_utc_time(end_time, field_name="end_time") async with conn.cursor() as cur: await cur.execute( "SELECT * FROM realtime.node_simulation WHERE time >= %s AND time <= %s", - (start_time, end_time), + (normalized_start_time, normalized_end_time), ) return await cur.fetchall() @@ -397,24 +399,9 @@ class RealtimeRepository: link_result_list: List of link simulation results result_start_time: Start time for the results (ISO format string) """ - # Convert result_start_time string to datetime if needed - if isinstance(result_start_time, str): - # 如果是ISO格式字符串,解析并转换为UTC+8 - if result_start_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat( - result_start_time.replace("Z", "+00:00") - ) - simulation_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - simulation_time = datetime.fromisoformat(result_start_time) - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) - else: - simulation_time = result_start_time - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) + simulation_time = parse_utc_time( + result_start_time, field_name="result_start_time" + ) # Prepare node data for batch insert node_data = [] @@ -475,24 +462,9 @@ class RealtimeRepository: link_result_list: List of link simulation results result_start_time: Start time for the results (ISO format string) """ - # Convert result_start_time string to datetime if needed - if isinstance(result_start_time, str): - # 如果是ISO格式字符串,解析并转换为UTC+8 - if result_start_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat( - result_start_time.replace("Z", "+00:00") - ) - simulation_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - simulation_time = datetime.fromisoformat(result_start_time) - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) - else: - simulation_time = result_start_time - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) + simulation_time = parse_utc_time( + result_start_time, field_name="result_start_time" + ) # Prepare node data for batch insert node_data = [] @@ -556,21 +528,7 @@ class RealtimeRepository: Returns: List of records matching the criteria """ - # Convert query_time string to datetime - if isinstance(query_time, str): - if query_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat(query_time.replace("Z", "+00:00")) - target_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - target_time = datetime.fromisoformat(query_time) - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) - else: - target_time = query_time - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) + target_time = parse_utc_time(query_time, field_name="query_time") # Create time range: query_time ± 1 second start_time = target_time - timedelta(seconds=1) @@ -614,21 +572,7 @@ class RealtimeRepository: Returns: List of records matching the criteria """ - # Convert query_time string to datetime - if isinstance(query_time, str): - if query_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat(query_time.replace("Z", "+00:00")) - target_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - target_time = datetime.fromisoformat(query_time) - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) - else: - target_time = query_time - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) + target_time = parse_utc_time(query_time, field_name="query_time") # Create time range: query_time ± 1 second start_time = target_time - timedelta(seconds=1) diff --git a/app/infra/db/timescaledb/repositories/scada.py b/app/infra/db/timescaledb/repositories/scada.py index bc8717f..d5a6348 100644 --- a/app/infra/db/timescaledb/repositories/scada.py +++ b/app/infra/db/timescaledb/repositories/scada.py @@ -89,12 +89,17 @@ class ScadaRepository: if field not in valid_fields: raise ValueError(f"Invalid field: {field}") - query = sql.SQL( + update_query = sql.SQL( "UPDATE scada.scada_data SET {} = %s WHERE time = %s AND device_id = %s" ).format(sql.Identifier(field)) + insert_query = sql.SQL( + "INSERT INTO scada.scada_data (time, device_id, {}) VALUES (%s, %s, %s)" + ).format(sql.Identifier(field)) async with conn.cursor() as cur: - await cur.execute(query, (value, time, device_id)) + await cur.execute(update_query, (value, time, device_id)) + if cur.rowcount == 0: + await cur.execute(insert_query, (time, device_id, value)) @staticmethod async def delete_scada_by_id_time_range( diff --git a/app/infra/db/timescaledb/repositories/scheme.py b/app/infra/db/timescaledb/repositories/scheme.py index bfa09ca..f0960b3 100644 --- a/app/infra/db/timescaledb/repositories/scheme.py +++ b/app/infra/db/timescaledb/repositories/scheme.py @@ -1,11 +1,9 @@ from typing import List, Any, Dict -from datetime import datetime, timedelta, timezone +from datetime import datetime, timedelta from collections import defaultdict from psycopg import AsyncConnection, Connection, sql import app.services.globals as globals - -# 定义UTC+8时区 -UTC_8 = timezone(timedelta(hours=8)) +from app.services.time_api import parse_utc_time class SchemeRepository: @@ -466,24 +464,9 @@ class SchemeRepository: link_result_list: List of link simulation results result_start_time: Start time for the results (ISO format string) """ - # Convert result_start_time string to datetime if needed - if isinstance(result_start_time, str): - # 如果是ISO格式字符串,解析并转换为UTC+8 - if result_start_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat( - result_start_time.replace("Z", "+00:00") - ) - simulation_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - simulation_time = datetime.fromisoformat(result_start_time) - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) - else: - simulation_time = result_start_time - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) + simulation_time = parse_utc_time( + result_start_time, field_name="result_start_time" + ) timestep_parts = globals.hydraulic_timestep.split(":") timestep = timedelta( @@ -564,24 +547,9 @@ class SchemeRepository: link_result_list: List of link simulation results result_start_time: Start time for the results (ISO format string) """ - # Convert result_start_time string to datetime if needed - if isinstance(result_start_time, str): - # 如果是ISO格式字符串,解析并转换为UTC+8 - if result_start_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat( - result_start_time.replace("Z", "+00:00") - ) - simulation_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - simulation_time = datetime.fromisoformat(result_start_time) - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) - else: - simulation_time = result_start_time - if simulation_time.tzinfo is None: - simulation_time = simulation_time.replace(tzinfo=UTC_8) + simulation_time = parse_utc_time( + result_start_time, field_name="result_start_time" + ) timestep_parts = globals.hydraulic_timestep.split(":") timestep = timedelta( @@ -664,21 +632,7 @@ class SchemeRepository: Returns: List of records matching the criteria """ - # Convert query_time string to datetime - if isinstance(query_time, str): - if query_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat(query_time.replace("Z", "+00:00")) - target_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - target_time = datetime.fromisoformat(query_time) - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) - else: - target_time = query_time - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) + target_time = parse_utc_time(query_time, field_name="query_time") # Create time range: query_time ± 1 second start_time = target_time - timedelta(seconds=1) @@ -727,21 +681,7 @@ class SchemeRepository: Returns: List of records matching the criteria """ - # Convert query_time string to datetime - if isinstance(query_time, str): - if query_time.endswith("Z"): - # UTC时间,转换为UTC+8 - utc_time = datetime.fromisoformat(query_time.replace("Z", "+00:00")) - target_time = utc_time.astimezone(UTC_8) - else: - # 假设已经是UTC+8时间 - target_time = datetime.fromisoformat(query_time) - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) - else: - target_time = query_time - if target_time.tzinfo is None: - target_time = target_time.replace(tzinfo=UTC_8) + target_time = parse_utc_time(query_time, field_name="query_time") # Create time range: query_time ± 1 second start_time = target_time - timedelta(seconds=1) diff --git a/app/services/time_api.py b/app/services/time_api.py index 85cc2f4..00904e2 100644 --- a/app/services/time_api.py +++ b/app/services/time_api.py @@ -1,5 +1,6 @@ -from datetime import datetime, timezone, timedelta -from dateutil import parser, tz +from datetime import date, datetime, time, timedelta, timezone + +from dateutil import parser, tz ''' 2025-02-09T15:45:00+00:00 采用的是 ISO 8601 国际标准日期时间格式,具体特点如下: @@ -13,57 +14,67 @@ from dateutil import parser, tz 2025-02-09T15:45:00+08:00 ''' -BG_TZ = tz.gettz('Asia/Shanghai') -UTC_TZ = tz.gettz('UTC') +BG_TZ = tz.gettz("Asia/Shanghai") +UTC_TZ = timezone.utc -def parse_utc_time(query_time: str) -> datetime: - ''' - 接受 任意格式的字符串,如果解析出来不带时区,则用 replace 添加 +00:00 时区 - 如果解析出来已经有时区,则用 astimezone 转换成UTC时间 - ''' +TIMEZONE_REQUIRED_MESSAGE = ( + "Datetime values must include an explicit timezone offset, for example " + "'2025-02-09T15:45:00Z' or '2025-02-09T23:45:00+08:00'." +) - # 解析时间字符串 - dt: datetime = parser.parse(query_time) + +def parse_aware_time(query_time: datetime | str, field_name: str = "datetime") -> datetime: + """ + 解析时间并确保结果带有时区信息。 + """ + dt = parser.parse(query_time) if isinstance(query_time, str) else query_time if dt.tzinfo is None: - dt = dt.replace(tzinfo=UTC_TZ) - else: - dt = dt.astimezone(UTC_TZ) - + raise ValueError(f"{field_name} is missing timezone information. {TIMEZONE_REQUIRED_MESSAGE}") return dt + + +def extract_date(value: date | datetime | str, field_name: str = "date") -> date: + """ + 提取日期部分,但保留调用方原始时区语义,不强制转换到 UTC。 + """ + if isinstance(value, date) and not isinstance(value, datetime): + return value + return parse_aware_time(value, field_name=field_name).date() + + +def utc_now() -> datetime: + """ + 返回带 UTC 时区的当前时间。 + """ + return datetime.now(UTC_TZ) + + +def parse_utc_time(query_time: datetime | str, field_name: str = "datetime") -> datetime: + ''' + 接受带时区的时间字符串/对象,并统一转换成 UTC 时间。 + ''' + return parse_aware_time(query_time, field_name=field_name).astimezone(UTC_TZ) -def parse_beijing_time(query_time: str) -> datetime: + +def parse_beijing_time(query_time: datetime | str, field_name: str = "datetime") -> datetime: ''' - 接受 任意格式的字符串,如果解析出来不带时区,则用 replace 添加 +08:00 时区 - 如果解析出来已经有时区,则用 astimezone 转换成北京时间 - - 也就是任意合法的时间字符串,最后都解析成 北京 时间 - + 接受带时区的时间字符串/对象,并统一转换成北京时间。 ''' - - # 解析时间字符串 - dt: datetime = parser.parse(query_time) - if dt.tzinfo is None: - dt = dt.replace(tzinfo=BG_TZ) - else: - dt = dt.astimezone(tz=BG_TZ) - - return dt + return parse_aware_time(query_time, field_name=field_name).astimezone(tz=BG_TZ) -def to_utc_time(dt: datetime) -> datetime: +def to_utc_time(dt: datetime | str, field_name: str = "datetime") -> datetime: ''' - 将一个北京时间的时间点,转换成utc + 将一个带时区的时间点,转换成 UTC。 ''' - utc_time = dt.astimezone(UTC_TZ) - return utc_time + return parse_aware_time(dt, field_name=field_name).astimezone(UTC_TZ) -def to_beijing_time(dt: datetime) -> datetime: +def to_beijing_time(dt: datetime | str, field_name: str = "datetime") -> datetime: ''' - 将一个 utc 的时间点,转换成北京时间 + 将一个带时区的时间点,转换成北京时间。 ''' - beijing_time = dt.astimezone(tz=BG_TZ) - return beijing_time + return parse_aware_time(dt, field_name=field_name).astimezone(tz=BG_TZ) def to_time_range(dt: datetime, delta: float) -> tuple[datetime, datetime]: @@ -83,7 +94,8 @@ def parse_beijing_date_range(query_date: str) -> tuple[datetime, datetime]: 将一个日期字符串,转换成 start/end 时间段,传进来的日期被认为是北京时间 日期字符串格式:YYYY-MM-DD ''' - start_time = parse_beijing_time(query_date) + target_date = date.fromisoformat(query_date) + start_time = datetime.combine(target_date, time.min, BG_TZ) end_time = start_time + timedelta(days=1) return (start_time, end_time) @@ -108,7 +120,7 @@ def get_date_from_time(time: str) -> str: ''' 将一个时间点,转换成日期 ''' - dt = parse_beijing_time(time) + dt = parse_beijing_time(time, field_name="time") return str(dt.date()) @@ -116,28 +128,27 @@ def is_today(query_date: str) -> bool: ''' 判断一个日期是否是今天 ''' - dt = parse_beijing_time(query_date) - return dt.date() == datetime.now().date() + dt = parse_beijing_time(query_date, field_name="query_date") + return dt.date() == datetime.now(BG_TZ).date() def is_yesterday(query_date: str) -> bool: ''' 判断一个日期是否是昨天 ''' - dt = parse_beijing_time(query_date) - return dt.date() == (datetime.now().date() - timedelta(days=1)) + dt = parse_beijing_time(query_date, field_name="query_date") + return dt.date() == (datetime.now(BG_TZ).date() - timedelta(days=1)) def is_tomorrow(query_date: str) -> bool: ''' 判断一个日期是否是明天 ''' - dt = parse_beijing_time(query_date) - return dt.date() == (datetime.now().date() + timedelta(days=1)) + dt = parse_beijing_time(query_date, field_name="query_date") + return dt.date() == (datetime.now(BG_TZ).date() + timedelta(days=1)) def is_today_or_future(query_date: str) -> bool: ''' 判断一个日期是否是今天或未来 ''' - dt = parse_beijing_time(query_date) - return dt.date() >= datetime.now().date() - + dt = parse_beijing_time(query_date, field_name="query_date") + return dt.date() >= datetime.now(BG_TZ).date()