fix(timeseries): align timezone and scada updates

This commit is contained in:
2026-06-11 10:27:51 +08:00
parent 24169bd277
commit 21929b44fc
9 changed files with 249 additions and 348 deletions
+5 -10
View File
@@ -8,8 +8,7 @@ from .dependencies import get_timescale_connection, get_postgres_connection
router = APIRouter() router = APIRouter()
@router.get("/composite/scada-simulation", summary="获取SCADA关联的模拟数据", @router.get("/composite/scada-simulation", summary="获取SCADA关联的模拟数据")
tags=["复合查询"])
async def get_scada_associated_simulation_data( async def get_scada_associated_simulation_data(
start_time: datetime = Query(..., description="查询开始时间"), start_time: datetime = Query(..., description="查询开始时间"),
end_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)) raise HTTPException(status_code=400, detail=str(e))
@router.get("/composite/element-simulation", summary="获取管网元素的模拟数据", @router.get("/composite/element-simulation", summary="获取管网元素的模拟数据")
tags=["复合查询"])
async def get_feature_simulation_data( async def get_feature_simulation_data(
start_time: datetime = Query(..., description="查询开始时间"), start_time: datetime = Query(..., description="查询开始时间"),
end_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)) raise HTTPException(status_code=400, detail=str(e))
@router.get("/composite/element-scada", summary="获取管网元素关联的SCADA监测数据", @router.get("/composite/element-scada", summary="获取管网元素关联的SCADA监测数据")
tags=["复合查询"])
async def get_element_associated_scada_data( async def get_element_associated_scada_data(
element_id: str = Query(..., description="管网元素ID(管道或节点)"), element_id: str = Query(..., description="管网元素ID(管道或节点)"),
start_time: datetime = Query(..., description="查询开始时间"), start_time: datetime = Query(..., description="查询开始时间"),
@@ -188,8 +185,7 @@ async def get_element_associated_scada_data(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.post("/composite/clean-scada", summary="清洗SCADA监测数据", @router.post("/composite/clean-scada", summary="清洗SCADA监测数据")
tags=["复合查询"])
async def clean_scada_data( async def clean_scada_data(
device_ids: str = Query(..., description="设备ID列表或 'all' 表示清洗所有设备"), device_ids: str = Query(..., description="设备ID列表或 'all' 表示清洗所有设备"),
start_time: datetime = Query(..., description="清洗数据的开始时间"), start_time: datetime = Query(..., description="清洗数据的开始时间"),
@@ -232,8 +228,7 @@ async def clean_scada_data(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.get("/composite/pipeline-health-prediction", summary="预测管道健康状况", @router.get("/composite/pipeline-health-prediction", summary="预测管道健康状况")
tags=["复合查询"])
async def predict_pipeline_health( async def predict_pipeline_health(
query_time: datetime = Query(..., description="查询时间"), query_time: datetime = Query(..., description="查询时间"),
network_name: str = Query(..., description="管网名称(或数据库名称)"), network_name: str = Query(..., description="管网名称(或数据库名称)"),
+54 -30
View File
@@ -8,9 +8,12 @@ from .dependencies import get_timescale_connection
router = APIRouter() 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( async def insert_realtime_links(
data: List[dict] = Body(..., description="管道数据列表,每项包含管道ID、时间戳等信息"), data: List[dict] = Body(..., description="管道数据列表,每项包含管道ID、时间戳等信息"),
conn: AsyncConnection = Depends(get_timescale_connection) conn: AsyncConnection = Depends(get_timescale_connection)
@@ -30,16 +33,21 @@ async def insert_realtime_links(
return {"message": f"Inserted {len(data)} records"} 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( async def get_realtime_links(
start_time: datetime = Query(..., description="查询开始时间"), start_time: datetime = Query(..., description=TIME_RANGE_START_DESC),
end_time: datetime = Query(..., description="查询结束时间"), end_time: datetime = Query(..., description=TIME_RANGE_END_DESC),
conn: AsyncConnection = Depends(get_timescale_connection), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
查询指定时间范围内的实时管道数据 查询指定时间范围内的实时管道数据
根据时间范围查询所有实时管道的监测值。 根据时间范围查询所有实时管道的监测值。传入时间必须显式包含时区,
可以直接使用 UTC+8,服务端会先统一转换为 UTC 再参与数据库查询。
Args: Args:
start_time: 查询开始时间 start_time: 查询开始时间
@@ -51,10 +59,14 @@ async def get_realtime_links(
return await RealtimeRepository.get_links_by_time_range(conn, start_time, end_time) 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( async def delete_realtime_links(
start_time: datetime = Query(..., description="删除开始时间"), start_time: datetime = Query(..., description=TIME_RANGE_START_DESC),
end_time: datetime = Query(..., description="删除结束时间"), end_time: datetime = Query(..., description=TIME_RANGE_END_DESC),
conn: AsyncConnection = Depends(get_timescale_connection), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
@@ -73,11 +85,10 @@ async def delete_realtime_links(
return {"message": "Deleted successfully"} return {"message": "Deleted successfully"}
@router.patch("/realtime/links/{link_id}/field", summary="更新实时管道字段", @router.patch("/realtime/links/{link_id}/field", summary="更新实时管道字段")
tags=["时间序列-实时数据"])
async def update_realtime_link_field( async def update_realtime_link_field(
link_id: str = Path(..., description="管道ID"), link_id: str = Path(..., description="管道ID"),
time: datetime = Query(..., description="更新数据的时间戳"), time: datetime = Query(..., description=f"更新记录的时间戳{TIME_WITH_TZ_DESC}"),
field: str = Query(..., description="要更新的字段名称"), field: str = Query(..., description="要更新的字段名称"),
value: float = Query(..., description="更新的字段值"), value: float = Query(..., description="更新的字段值"),
conn: AsyncConnection = Depends(get_timescale_connection), conn: AsyncConnection = Depends(get_timescale_connection),
@@ -106,8 +117,7 @@ async def update_realtime_link_field(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.post("/realtime/nodes/batch", status_code=201, summary="批量插入实时节点数据", @router.post("/realtime/nodes/batch", status_code=201, summary="批量插入实时节点数据")
tags=["时间序列-实时数据"])
async def insert_realtime_nodes( async def insert_realtime_nodes(
data: List[dict] = Body(..., description="节点数据列表,每项包含节点ID、时间戳等信息"), data: List[dict] = Body(..., description="节点数据列表,每项包含节点ID、时间戳等信息"),
conn: AsyncConnection = Depends(get_timescale_connection) conn: AsyncConnection = Depends(get_timescale_connection)
@@ -127,16 +137,21 @@ async def insert_realtime_nodes(
return {"message": f"Inserted {len(data)} records"} 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( async def get_realtime_nodes(
start_time: datetime = Query(..., description="查询开始时间"), start_time: datetime = Query(..., description=TIME_RANGE_START_DESC),
end_time: datetime = Query(..., description="查询结束时间"), end_time: datetime = Query(..., description=TIME_RANGE_END_DESC),
conn: AsyncConnection = Depends(get_timescale_connection), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
查询指定时间范围内的实时节点数据 查询指定时间范围内的实时节点数据
根据时间范围查询所有实时节点的监测值。 根据时间范围查询所有实时节点的监测值。传入时间必须显式包含时区,
可以直接使用 UTC+8,服务端会先统一转换为 UTC 再参与数据库查询。
Args: Args:
start_time: 查询开始时间 start_time: 查询开始时间
@@ -148,10 +163,14 @@ async def get_realtime_nodes(
return await RealtimeRepository.get_nodes_by_time_range(conn, start_time, end_time) 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( async def delete_realtime_nodes(
start_time: datetime = Query(..., description="删除开始时间"), start_time: datetime = Query(..., description=TIME_RANGE_START_DESC),
end_time: datetime = Query(..., description="删除结束时间"), end_time: datetime = Query(..., description=TIME_RANGE_END_DESC),
conn: AsyncConnection = Depends(get_timescale_connection), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
@@ -172,12 +191,11 @@ async def delete_realtime_nodes(
@router.post("/realtime/simulation/store", status_code=201, summary="存储实时模拟结果", @router.post("/realtime/simulation/store", status_code=201, summary="存储实时模拟结果")
tags=["时间序列-实时数据"])
async def store_realtime_simulation_result( async def store_realtime_simulation_result(
node_result_list: List[dict] = Body(..., description="节点模拟结果列表"), node_result_list: List[dict] = Body(..., description="节点模拟结果列表"),
link_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), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
@@ -199,10 +217,13 @@ async def store_realtime_simulation_result(
return {"message": "Simulation results stored successfully"} return {"message": "Simulation results stored successfully"}
@router.get("/realtime/query/by-time-property", summary="按时间和属性查询实时数据", @router.get(
tags=["时间序列-实时数据"]) "/realtime/query/by-time-property",
summary="按时间和属性查询实时数据",
description="查询指定时间点的实时属性值。query_time 必须显式带时区;允许传 UTC+8,服务端会先归一化为 UTC 再执行查询。",
)
async def query_realtime_records_by_time_property( 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(节点)"), type: str = Query(..., description="数据类型,pipe(管道)或 junction(节点)"),
property: str = Query(..., description="要查询的属性名称"), property: str = Query(..., description="要查询的属性名称"),
conn: AsyncConnection = Depends(get_timescale_connection), 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)) raise HTTPException(status_code=400, detail=str(e))
@router.get("/realtime/query/by-id-time", summary="按ID和时间查询实时模拟数据", @router.get(
tags=["时间序列-实时数据"]) "/realtime/query/by-id-time",
summary="按ID和时间查询实时模拟数据",
description="查询指定元素在某一时间点的实时模拟结果。query_time 必须显式带时区;允许传 UTC+8,服务端会先归一化为 UTC 再执行查询。",
)
async def query_realtime_simulation_by_id_time( async def query_realtime_simulation_by_id_time(
id: str = Query(..., description="元素ID(管道ID或节点ID"), id: str = Query(..., description="元素ID(管道ID或节点ID"),
type: str = Query(..., description="元素类型,pipe(管道)或 junction(节点)"), 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), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
+14 -13
View File
@@ -9,11 +9,10 @@ from .dependencies import get_timescale_connection
router = APIRouter() router = APIRouter()
@router.post("/scada/batch", status_code=201, summary="批量插入SCADA监测数据", @router.post("/scada/batch", status_code=201, summary="批量插入SCADA监测数据")
tags=["时间序列-监测数据"])
async def insert_scada_data( async def insert_scada_data(
data: List[dict] = Body(..., description="SCADA设备监测数据列表"), data: List[dict] = Body(..., description="SCADA设备监测数据列表"),
conn: AsyncConnection = Depends(get_timescale_connection) conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
批量插入SCADA监测数据 批量插入SCADA监测数据
@@ -30,12 +29,13 @@ async def insert_scada_data(
return {"message": f"Inserted {len(data)} records"} return {"message": f"Inserted {len(data)} records"}
@router.get("/scada/by-ids-time-range", summary="按设备ID和时间范围查询SCADA数据", @router.get("/scada/by-ids-time-range", summary="按设备ID和时间范围查询SCADA数据")
tags=["时间序列-监测数据"])
async def get_scada_by_ids_time_range( async def get_scada_by_ids_time_range(
start_time: datetime = Query(..., description="查询开始时间"), start_time: datetime = Query(..., description="查询开始时间"),
end_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), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
@@ -59,13 +59,16 @@ async def get_scada_by_ids_time_range(
) )
@router.get("/scada/by-ids-field-time-range", summary="按设备ID、字段和时间范围查询SCADA数据", @router.get(
tags=["时间序列-监测数据"]) "/scada/by-ids-field-time-range", summary="按设备ID、字段和时间范围查询SCADA数据"
)
async def get_scada_field_by_ids_time_range( async def get_scada_field_by_ids_time_range(
start_time: datetime = Query(..., description="查询开始时间"), start_time: datetime = Query(..., description="查询开始时间"),
end_time: datetime = Query(..., description="查询结束时间"), end_time: datetime = Query(..., description="查询结束时间"),
field: str = 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), conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
@@ -98,8 +101,7 @@ async def get_scada_field_by_ids_time_range(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.patch("/scada/{device_id}/field", summary="更新SCADA设备字段", @router.patch("/scada/{device_id}/field", summary="更新SCADA设备字段")
tags=["时间序列-监测数据"])
async def update_scada_field( async def update_scada_field(
device_id: str = Path(..., description="设备ID"), device_id: str = Path(..., description="设备ID"),
time: datetime = Query(..., description="更新数据的时间戳"), time: datetime = Query(..., description="更新数据的时间戳"),
@@ -131,8 +133,7 @@ async def update_scada_field(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.delete("/scada/by-id-time-range", summary="按设备ID和时间范围删除SCADA数据", @router.delete("/scada/by-id-time-range", summary="按设备ID和时间范围删除SCADA数据")
tags=["时间序列-监测数据"])
async def delete_scada_data( async def delete_scada_data(
device_id: str = Query(..., description="设备ID"), device_id: str = Query(..., description="设备ID"),
start_time: datetime = Query(..., description="删除开始时间"), start_time: datetime = Query(..., description="删除开始时间"),
+16 -23
View File
@@ -9,11 +9,10 @@ from .dependencies import get_timescale_connection
router = APIRouter() router = APIRouter()
@router.post("/scheme/links/batch", status_code=201, summary="批量插入方案管道数据", @router.post("/scheme/links/batch", status_code=201, summary="批量插入方案管道数据")
tags=["时间序列-方案数据"])
async def insert_scheme_links( async def insert_scheme_links(
data: List[dict] = Body(..., description="方案管道数据列表"), data: List[dict] = Body(..., description="方案管道数据列表"),
conn: AsyncConnection = Depends(get_timescale_connection) conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
批量插入方案管道数据 批量插入方案管道数据
@@ -30,7 +29,7 @@ async def insert_scheme_links(
return {"message": f"Inserted {len(data)} records"} return {"message": f"Inserted {len(data)} records"}
@router.get("/scheme/links", summary="查询方案管道数据", tags=["时间序列-方案数据"]) @router.get("/scheme/links", summary="查询方案管道数据")
async def get_scheme_links( async def get_scheme_links(
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
scheme_name: str = Query(..., description="方案名称"), scheme_name: str = Query(..., description="方案名称"),
@@ -57,8 +56,7 @@ async def get_scheme_links(
) )
@router.get("/scheme/links/{link_id}/field", summary="查询方案管道字段数据", @router.get("/scheme/links/{link_id}/field", summary="查询方案管道字段数据")
tags=["时间序列-方案数据"])
async def get_scheme_link_field( async def get_scheme_link_field(
link_id: str = Path(..., description="管道ID"), link_id: str = Path(..., description="管道ID"),
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
@@ -95,8 +93,7 @@ async def get_scheme_link_field(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.patch("/scheme/links/{link_id}/field", summary="更新方案管道字段", @router.patch("/scheme/links/{link_id}/field", summary="更新方案管道字段")
tags=["时间序列-方案数据"])
async def update_scheme_link_field( async def update_scheme_link_field(
link_id: str = Path(..., description="管道ID"), link_id: str = Path(..., description="管道ID"),
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
@@ -134,7 +131,7 @@ async def update_scheme_link_field(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.delete("/scheme/links", summary="删除方案管道数据", tags=["时间序列-方案数据"]) @router.delete("/scheme/links", summary="删除方案管道数据")
async def delete_scheme_links( async def delete_scheme_links(
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
scheme_name: str = Query(..., description="方案名称"), scheme_name: str = Query(..., description="方案名称"),
@@ -162,11 +159,10 @@ async def delete_scheme_links(
return {"message": "Deleted successfully"} return {"message": "Deleted successfully"}
@router.post("/scheme/nodes/batch", status_code=201, summary="批量插入方案节点数据", @router.post("/scheme/nodes/batch", status_code=201, summary="批量插入方案节点数据")
tags=["时间序列-方案数据"])
async def insert_scheme_nodes( async def insert_scheme_nodes(
data: List[dict] = Body(..., description="方案节点数据列表"), data: List[dict] = Body(..., description="方案节点数据列表"),
conn: AsyncConnection = Depends(get_timescale_connection) conn: AsyncConnection = Depends(get_timescale_connection),
): ):
""" """
批量插入方案节点数据 批量插入方案节点数据
@@ -183,8 +179,7 @@ async def insert_scheme_nodes(
return {"message": f"Inserted {len(data)} records"} return {"message": f"Inserted {len(data)} records"}
@router.get("/scheme/nodes/{node_id}/field", summary="查询方案节点字段数据", @router.get("/scheme/nodes/{node_id}/field", summary="查询方案节点字段数据")
tags=["时间序列-方案数据"])
async def get_scheme_node_field( async def get_scheme_node_field(
node_id: str = Path(..., description="节点ID"), node_id: str = Path(..., description="节点ID"),
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
@@ -221,8 +216,7 @@ async def get_scheme_node_field(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.patch("/scheme/nodes/{node_id}/field", summary="更新方案节点字段", @router.patch("/scheme/nodes/{node_id}/field", summary="更新方案节点字段")
tags=["时间序列-方案数据"])
async def update_scheme_node_field( async def update_scheme_node_field(
node_id: str = Path(..., description="节点ID"), node_id: str = Path(..., description="节点ID"),
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
@@ -260,7 +254,7 @@ async def update_scheme_node_field(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.delete("/scheme/nodes", summary="删除方案节点数据", tags=["时间序列-方案数据"]) @router.delete("/scheme/nodes", summary="删除方案节点数据")
async def delete_scheme_nodes( async def delete_scheme_nodes(
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
scheme_name: str = Query(..., description="方案名称"), scheme_name: str = Query(..., description="方案名称"),
@@ -288,8 +282,7 @@ async def delete_scheme_nodes(
return {"message": "Deleted successfully"} return {"message": "Deleted successfully"}
@router.post("/scheme/simulation/store", status_code=201, summary="存储方案模拟结果", @router.post("/scheme/simulation/store", status_code=201, summary="存储方案模拟结果")
tags=["时间序列-方案数据"])
async def store_scheme_simulation_result( async def store_scheme_simulation_result(
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
scheme_name: str = Query(..., description="方案名称"), scheme_name: str = Query(..., description="方案名称"),
@@ -324,8 +317,9 @@ async def store_scheme_simulation_result(
return {"message": "Scheme simulation results stored successfully"} return {"message": "Scheme simulation results stored successfully"}
@router.get("/scheme/query/by-scheme-time-property", summary="按方案、时间和属性查询数据", @router.get(
tags=["时间序列-方案数据"]) "/scheme/query/by-scheme-time-property", summary="按方案、时间和属性查询数据"
)
async def query_scheme_records_by_scheme_time_property( async def query_scheme_records_by_scheme_time_property(
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
scheme_name: str = Query(..., description="方案名称"), scheme_name: str = Query(..., description="方案名称"),
@@ -361,8 +355,7 @@ async def query_scheme_records_by_scheme_time_property(
raise HTTPException(status_code=400, detail=str(e)) raise HTTPException(status_code=400, detail=str(e))
@router.get("/scheme/query/by-id-time", summary="按ID和时间查询方案模拟数据", @router.get("/scheme/query/by-id-time", summary="按ID和时间查询方案模拟数据")
tags=["时间序列-方案数据"])
async def query_scheme_simulation_by_id_time( async def query_scheme_simulation_by_id_time(
scheme_type: str = Query(..., description="方案类型"), scheme_type: str = Query(..., description="方案类型"),
scheme_name: str = Query(..., description="方案名称"), scheme_name: str = Query(..., description="方案名称"),
+8 -20
View File
@@ -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.scheme import SchemeRepository
from app.infra.db.timescaledb.repositories.realtime import RealtimeRepository from app.infra.db.timescaledb.repositories.realtime import RealtimeRepository
from app.infra.db.timescaledb.repositories.scada import ScadaRepository from app.infra.db.timescaledb.repositories.scada import ScadaRepository
from app.services.time_api import parse_utc_time
class InternalStorage: class InternalStorage:
@@ -89,10 +90,9 @@ class InternalQueries:
) -> dict: ) -> dict:
"""查询指定时间点的 SCADA 数据""" """查询指定时间点的 SCADA 数据"""
# 解析时间,假设是北京时间 target_time = parse_utc_time(query_time, field_name="query_time")
beijing_time = datetime.fromisoformat(query_time) start_time = target_time - timedelta(seconds=1)
start_time = beijing_time - timedelta(seconds=1) end_time = target_time + timedelta(seconds=1)
end_time = beijing_time + timedelta(seconds=1)
for attempt in range(max_retries): for attempt in range(max_retries):
try: try:
@@ -132,14 +132,8 @@ class InternalQueries:
max_retries: int = 3, max_retries: int = 3,
) -> dict[str, list[dict]]: ) -> dict[str, list[dict]]:
"""查询指定时间窗的 SCADA 数据,返回 {device_id: [{time, value}, ...]}。""" """查询指定时间窗的 SCADA 数据,返回 {device_id: [{time, value}, ...]}。"""
start_dt = ( start_dt = parse_utc_time(start_time, field_name="start_time")
datetime.fromisoformat(start_time) end_dt = parse_utc_time(end_time, field_name="end_time")
if isinstance(start_time, str)
else start_time
)
end_dt = (
datetime.fromisoformat(end_time) if isinstance(end_time, str) else end_time
)
for attempt in range(max_retries): for attempt in range(max_retries):
try: try:
@@ -238,14 +232,8 @@ class InternalQueries:
if not element_ids: if not element_ids:
return {} return {}
start_dt = ( start_dt = parse_utc_time(start_time, field_name="start_time")
datetime.fromisoformat(start_time) end_dt = parse_utc_time(end_time, field_name="end_time")
if isinstance(start_time, str)
else start_time
)
end_dt = (
datetime.fromisoformat(end_time) if isinstance(end_time, str) else end_time
)
table_name, valid_fields = InternalQueries._resolve_simulation_table(element_type) table_name, valid_fields = InternalQueries._resolve_simulation_table(element_type)
if field not in valid_fields: if field not in valid_fields:
raise ValueError(f"Invalid field for {element_type}: {field}") raise ValueError(f"Invalid field for {element_type}: {field}")
@@ -1,10 +1,8 @@
from typing import List, Any, Dict from typing import List, Any, Dict
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta
from collections import defaultdict from collections import defaultdict
from psycopg import AsyncConnection, Connection, sql from psycopg import AsyncConnection, Connection, sql
from app.services.time_api import parse_utc_time
# 定义UTC+8时区
UTC_8 = timezone(timedelta(hours=8))
class RealtimeRepository: class RealtimeRepository:
@@ -102,10 +100,12 @@ class RealtimeRepository:
async def get_links_by_time_range( async def get_links_by_time_range(
conn: AsyncConnection, start_time: datetime, end_time: datetime conn: AsyncConnection, start_time: datetime, end_time: datetime
) -> List[dict]: ) -> 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: async with conn.cursor() as cur:
await cur.execute( await cur.execute(
"SELECT * FROM realtime.link_simulation WHERE time >= %s AND time <= %s", "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() return await cur.fetchall()
@@ -298,10 +298,12 @@ class RealtimeRepository:
async def get_nodes_by_time_range( async def get_nodes_by_time_range(
conn: AsyncConnection, start_time: datetime, end_time: datetime conn: AsyncConnection, start_time: datetime, end_time: datetime
) -> List[dict]: ) -> 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: async with conn.cursor() as cur:
await cur.execute( await cur.execute(
"SELECT * FROM realtime.node_simulation WHERE time >= %s AND time <= %s", "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() return await cur.fetchall()
@@ -397,24 +399,9 @@ class RealtimeRepository:
link_result_list: List of link simulation results link_result_list: List of link simulation results
result_start_time: Start time for the results (ISO format string) result_start_time: Start time for the results (ISO format string)
""" """
# Convert result_start_time string to datetime if needed simulation_time = parse_utc_time(
if isinstance(result_start_time, str): result_start_time, field_name="result_start_time"
# 如果是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)
# Prepare node data for batch insert # Prepare node data for batch insert
node_data = [] node_data = []
@@ -475,24 +462,9 @@ class RealtimeRepository:
link_result_list: List of link simulation results link_result_list: List of link simulation results
result_start_time: Start time for the results (ISO format string) result_start_time: Start time for the results (ISO format string)
""" """
# Convert result_start_time string to datetime if needed simulation_time = parse_utc_time(
if isinstance(result_start_time, str): result_start_time, field_name="result_start_time"
# 如果是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)
# Prepare node data for batch insert # Prepare node data for batch insert
node_data = [] node_data = []
@@ -556,21 +528,7 @@ class RealtimeRepository:
Returns: Returns:
List of records matching the criteria List of records matching the criteria
""" """
# Convert query_time string to datetime target_time = parse_utc_time(query_time, field_name="query_time")
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)
# Create time range: query_time ± 1 second # Create time range: query_time ± 1 second
start_time = target_time - timedelta(seconds=1) start_time = target_time - timedelta(seconds=1)
@@ -614,21 +572,7 @@ class RealtimeRepository:
Returns: Returns:
List of records matching the criteria List of records matching the criteria
""" """
# Convert query_time string to datetime target_time = parse_utc_time(query_time, field_name="query_time")
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)
# Create time range: query_time ± 1 second # Create time range: query_time ± 1 second
start_time = target_time - timedelta(seconds=1) start_time = target_time - timedelta(seconds=1)
@@ -89,12 +89,17 @@ class ScadaRepository:
if field not in valid_fields: if field not in valid_fields:
raise ValueError(f"Invalid field: {field}") 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" "UPDATE scada.scada_data SET {} = %s WHERE time = %s AND device_id = %s"
).format(sql.Identifier(field)) ).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: 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 @staticmethod
async def delete_scada_by_id_time_range( async def delete_scada_by_id_time_range(
@@ -1,11 +1,9 @@
from typing import List, Any, Dict from typing import List, Any, Dict
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta
from collections import defaultdict from collections import defaultdict
from psycopg import AsyncConnection, Connection, sql from psycopg import AsyncConnection, Connection, sql
import app.services.globals as globals import app.services.globals as globals
from app.services.time_api import parse_utc_time
# 定义UTC+8时区
UTC_8 = timezone(timedelta(hours=8))
class SchemeRepository: class SchemeRepository:
@@ -466,24 +464,9 @@ class SchemeRepository:
link_result_list: List of link simulation results link_result_list: List of link simulation results
result_start_time: Start time for the results (ISO format string) result_start_time: Start time for the results (ISO format string)
""" """
# Convert result_start_time string to datetime if needed simulation_time = parse_utc_time(
if isinstance(result_start_time, str): result_start_time, field_name="result_start_time"
# 如果是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)
timestep_parts = globals.hydraulic_timestep.split(":") timestep_parts = globals.hydraulic_timestep.split(":")
timestep = timedelta( timestep = timedelta(
@@ -564,24 +547,9 @@ class SchemeRepository:
link_result_list: List of link simulation results link_result_list: List of link simulation results
result_start_time: Start time for the results (ISO format string) result_start_time: Start time for the results (ISO format string)
""" """
# Convert result_start_time string to datetime if needed simulation_time = parse_utc_time(
if isinstance(result_start_time, str): result_start_time, field_name="result_start_time"
# 如果是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)
timestep_parts = globals.hydraulic_timestep.split(":") timestep_parts = globals.hydraulic_timestep.split(":")
timestep = timedelta( timestep = timedelta(
@@ -664,21 +632,7 @@ class SchemeRepository:
Returns: Returns:
List of records matching the criteria List of records matching the criteria
""" """
# Convert query_time string to datetime target_time = parse_utc_time(query_time, field_name="query_time")
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)
# Create time range: query_time ± 1 second # Create time range: query_time ± 1 second
start_time = target_time - timedelta(seconds=1) start_time = target_time - timedelta(seconds=1)
@@ -727,21 +681,7 @@ class SchemeRepository:
Returns: Returns:
List of records matching the criteria List of records matching the criteria
""" """
# Convert query_time string to datetime target_time = parse_utc_time(query_time, field_name="query_time")
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)
# Create time range: query_time ± 1 second # Create time range: query_time ± 1 second
start_time = target_time - timedelta(seconds=1) start_time = target_time - timedelta(seconds=1)
+64 -53
View File
@@ -1,4 +1,5 @@
from datetime import datetime, timezone, timedelta from datetime import date, datetime, time, timedelta, timezone
from dateutil import parser, tz from dateutil import parser, tz
''' '''
@@ -13,57 +14,67 @@ from dateutil import parser, tz
2025-02-09T15:45:00+08:00 2025-02-09T15:45:00+08:00
''' '''
BG_TZ = tz.gettz('Asia/Shanghai') BG_TZ = tz.gettz("Asia/Shanghai")
UTC_TZ = tz.gettz('UTC') UTC_TZ = timezone.utc
def parse_utc_time(query_time: str) -> datetime: TIMEZONE_REQUIRED_MESSAGE = (
''' "Datetime values must include an explicit timezone offset, for example "
接受 任意格式的字符串,如果解析出来不带时区,则用 replace 添加 +00:00 时区 "'2025-02-09T15:45:00Z' or '2025-02-09T23:45:00+08:00'."
如果解析出来已经有时区,则用 astimezone 转换成UTC时间 )
'''
# 解析时间字符串
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: if dt.tzinfo is None:
dt = dt.replace(tzinfo=UTC_TZ) raise ValueError(f"{field_name} is missing timezone information. {TIMEZONE_REQUIRED_MESSAGE}")
else:
dt = dt.astimezone(UTC_TZ)
return dt
def parse_beijing_time(query_time: str) -> 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 dt
def to_utc_time(dt: datetime) -> datetime: def extract_date(value: date | datetime | str, field_name: str = "date") -> date:
''' """
将一个北京时间的时间点,转换成utc 提取日期部分,但保留调用方原始时区语义,不强制转换到 UTC。
''' """
utc_time = dt.astimezone(UTC_TZ) if isinstance(value, date) and not isinstance(value, datetime):
return utc_time return value
return parse_aware_time(value, field_name=field_name).date()
def to_beijing_time(dt: datetime) -> datetime: def utc_now() -> datetime:
"""
返回带 UTC 时区的当前时间。
"""
return datetime.now(UTC_TZ)
def parse_utc_time(query_time: datetime | str, field_name: str = "datetime") -> datetime:
''' '''
将一个 utc 的时间点,转换成北京时间 接受带时区的时间字符串/对象,并统一转换成 UTC 时间
''' '''
beijing_time = dt.astimezone(tz=BG_TZ) return parse_aware_time(query_time, field_name=field_name).astimezone(UTC_TZ)
return beijing_time
def parse_beijing_time(query_time: datetime | str, field_name: str = "datetime") -> datetime:
'''
接受带时区的时间字符串/对象,并统一转换成北京时间。
'''
return parse_aware_time(query_time, field_name=field_name).astimezone(tz=BG_TZ)
def to_utc_time(dt: datetime | str, field_name: str = "datetime") -> datetime:
'''
将一个带时区的时间点,转换成 UTC。
'''
return parse_aware_time(dt, field_name=field_name).astimezone(UTC_TZ)
def to_beijing_time(dt: datetime | str, field_name: str = "datetime") -> datetime:
'''
将一个带时区的时间点,转换成北京时间。
'''
return parse_aware_time(dt, field_name=field_name).astimezone(tz=BG_TZ)
def to_time_range(dt: datetime, delta: float) -> tuple[datetime, datetime]: 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 时间段,传进来的日期被认为是北京时间 将一个日期字符串,转换成 start/end 时间段,传进来的日期被认为是北京时间
日期字符串格式:YYYY-MM-DD 日期字符串格式: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) end_time = start_time + timedelta(days=1)
return (start_time, end_time) 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()) return str(dt.date())
@@ -116,28 +128,27 @@ def is_today(query_date: str) -> bool:
''' '''
判断一个日期是否是今天 判断一个日期是否是今天
''' '''
dt = parse_beijing_time(query_date) dt = parse_beijing_time(query_date, field_name="query_date")
return dt.date() == datetime.now().date() return dt.date() == datetime.now(BG_TZ).date()
def is_yesterday(query_date: str) -> bool: def is_yesterday(query_date: str) -> bool:
''' '''
判断一个日期是否是昨天 判断一个日期是否是昨天
''' '''
dt = parse_beijing_time(query_date) dt = parse_beijing_time(query_date, field_name="query_date")
return dt.date() == (datetime.now().date() - timedelta(days=1)) return dt.date() == (datetime.now(BG_TZ).date() - timedelta(days=1))
def is_tomorrow(query_date: str) -> bool: def is_tomorrow(query_date: str) -> bool:
''' '''
判断一个日期是否是明天 判断一个日期是否是明天
''' '''
dt = parse_beijing_time(query_date) dt = parse_beijing_time(query_date, field_name="query_date")
return dt.date() == (datetime.now().date() + timedelta(days=1)) return dt.date() == (datetime.now(BG_TZ).date() + timedelta(days=1))
def is_today_or_future(query_date: str) -> bool: def is_today_or_future(query_date: str) -> bool:
''' '''
判断一个日期是否是今天或未来 判断一个日期是否是今天或未来
''' '''
dt = parse_beijing_time(query_date) dt = parse_beijing_time(query_date, field_name="query_date")
return dt.date() >= datetime.now().date() return dt.date() >= datetime.now(BG_TZ).date()