Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions pandahub/api/routers/timeseries.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ class GetTimeSeriesModel(BaseModel):
timestamp_range: Optional[tuple] = None
exclude_timestamp_range: Optional[tuple] = None
collection_name: Optional[str] = "timeseries"
as_utc: Optional[bool] = False


@router.post("/get_timeseries_from_db")
Expand All @@ -39,6 +40,7 @@ class MultiGetTimeSeriesModel(BaseModel):
timestamp_range: Optional[tuple] = None
exclude_timestamp_range: Optional[tuple] = None
collection_name: Optional[str] = "timeseries"
as_utc: Optional[bool] = False


@router.post("/multi_get_timeseries_from_db")
Expand Down
9 changes: 8 additions & 1 deletion pandahub/lib/PandaHub.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from inspect import signature, _empty
from collections.abc import Callable
from typing import Optional, Union, TypeVar
from datetime import timezone

import numpy as np
import pandas as pd
Expand Down Expand Up @@ -2214,6 +2215,7 @@ def get_timeseries_from_db(
collection_name="timeseries",
include_metadata=False,
project_id=None,
as_utc=False,
**kwargs,
):
"""
Expand Down Expand Up @@ -2350,6 +2352,8 @@ def get_timeseries_from_db(
timeseries_data = pd.Series(
data["values"], index=data["timestamps"], dtype="float64"
)
if as_utc:
timeseries_data = timeseries_data.tz_localize("UTC")
elif ts_format == "array":
timeseries_data = data["timeseries_data"]
if include_metadata:
Expand Down Expand Up @@ -2487,6 +2491,7 @@ def multi_get_timeseries_from_db(
global_database=False,
collection_name="timeseries",
project_id=None,
as_utc=False,
**kwargs,
):
if project_id:
Expand Down Expand Up @@ -2649,12 +2654,14 @@ def multi_get_timeseries_from_db(
continue
data = ts["timeseries_data"]
if compressed_ts_data:
timeseries_data = decompress_timeseries_data(data, ts_format)
timeseries_data = decompress_timeseries_data(data, ts_format, as_utc)
ts["timeseries_data"] = timeseries_data
else:
if ts_format == "timestamp_value":
timeseries_data = pd.DataFrame(ts["timeseries_data"])
timeseries_data.set_index("timestamp", inplace=True)
if as_utc:
timeseries_data = timeseries_data.tz_localize("UTC")
timeseries_data.index.name = None
ts["timeseries_data"] = timeseries_data.value
if include_metadata:
Expand Down
8 changes: 6 additions & 2 deletions pandahub/lib/database_toolbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
logger = logging.getLogger(__name__)
from pandapower.io_utils import PPJSONEncoder
from packaging import version
from datetime import timezone


def get_document_hash(task):
Expand Down Expand Up @@ -125,12 +126,15 @@ def compress_timeseries_data(timeseries_data, ts_format):
cname="zlib")


def decompress_timeseries_data(timeseries_data, ts_format, num_timestamps):
def decompress_timeseries_data(timeseries_data, ts_format, num_timestamps, as_utc=False):
if ts_format == "timestamp_value":
data = np.frombuffer(blosc.decompress(timeseries_data),
dtype=np.float64).reshape((num_timestamps, 2),
order="F")
return pd.Series(data[:,1], index=pd.to_datetime(data[:,0]))
data = pd.Series(data[:, 1], index=pd.to_datetime(data[:, 0]))
if as_utc:
data = data.tz_localize("UTC")
return data
elif ts_format == "array":
return np.frombuffer(blosc.decompress(timeseries_data),
dtype=np.float64)
Expand Down