Skip to content
Merged
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
7 changes: 5 additions & 2 deletions environment.yml
Original file line number Diff line number Diff line change
Expand Up @@ -18,14 +18,15 @@ channels:
- conda-forge
dependencies:
- python>=3.10,<3.13
- setuptools>=68.0.0
- importlib-metadata>=7.0.0,<9.0.0
- jinja2>=3.1.6,<4.0.0
- pytest==7.4.0
- pytest-mock==3.11.1
- pytest-cov==4.1.0
- pylint==2.17.4
- pip>=23.1.2
- turbodbc==4.11.0
- turbodbc>=4.5.0
- numpy>=1.23.4,<2.0.0
- oauthlib>=3.2.2,<4.0.0
- cryptography>=38.0.3
Expand Down Expand Up @@ -74,6 +75,8 @@ dependencies:
- pmdarima>=2.0.4
- scikit-learn>=1.3.0,<1.6.0
- pip:
# Ensure setuptools is installed first for pkg_resources needed by turbodbc
- setuptools>=68.0.0
# protobuf installed via pip to avoid libabseil conflicts with conda libarrow
- protobuf>=5.29.0,<5.30.0
# pyarrow constraint must match conda version to avoid upgrades
Expand All @@ -91,4 +94,4 @@ dependencies:
- web3>=7.7.0,<8.0.0
- eth-typing>=5.0.1,<6.0.0
- pandas>=2.0.1,<2.3.0
- moto[s3]>=5.0.16,<6.0.0
- moto[s3]>=5.0.16,<6.0.0
2 changes: 2 additions & 0 deletions setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@
long_description = (here / "PYPI-README.md").read_text()

INSTALL_REQUIRES = [
"setuptools>=68.0.0",
# "turbodbc>=4.5.0",
"databricks-sql-connector>=3.6.0,<3.7.0",
"pyarrow>=14.0.1,<17.0.0",
"azure-core>=1.38.0,<2.0.0",
Expand Down
12 changes: 5 additions & 7 deletions src/api/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -22,21 +22,20 @@ ENV AzureWebJobsScriptRoot=/home/site/wwwroot \
COPY src/api/requirements.txt /

RUN rm -rf /var/lib/apt/lists/partial \
apt-get clean \
&& apt-get clean \
&& apt-get update -o Acquire::CompressionTypes::Order::=gz \
&& apt-get --no-install-recommends install -y odbcinst1debian2 libodbc1 odbcinst unixodbc unixodbc-dev libsasl2-dev libsasl2-modules-gssapi-mit libboost-all-dev \
&& apt-get --no-install-recommends install -y ca-certificates curl python3-pip python3-dev python3-setuptools python3-wheel gcc g++ \
&& apt-get --no-install-recommends install -y zip unzip wget \
&& apt-get --no-install-recommends install -y zip unzip wget \
&& wget --secure-protocol=TLSv1_2 --max-redirect=0 https://databricks-bi-artifacts.s3.us-east-2.amazonaws.com/simbaspark-drivers/odbc/2.7.7/SimbaSparkODBC-2.7.7.1016-Debian-64bit.zip -P /odbc/ \
&& unzip /odbc/SimbaSparkODBC-2.7.7.1016-Debian-64bit.zip -d /odbc \
&& dpkg -i /odbc/simbaspark_2.7.7.1016-2_amd64.deb \
&& pip install --no-cache-dir pyarrow==14.0.2 \
&& pip install --no-cache-dir numpy==1.26.4 \
&& python -c "import pyarrow; pyarrow.create_library_symlinks()" \
&& CFLAGS="-D_GLIBCXX_USE_CXX11_ABI=0" pip install --no-cache-dir -r /requirements.txt \
&& rm -rf /var/lib/apt/lists/* /var/cache/apt/archives/* \
&& rm /odbc -r

&& rm -rf /var/lib/apt/lists/* /var/cache/apt/archives/* \
&& rm -rf /odbc
RUN echo '[ODBC Drivers]' > /etc/odbcinst.ini \
&& echo 'Simba Spark ODBC Driver = Installed' >> /etc/odbcinst.ini \
&& echo '[Simba Spark ODBC Driver]' >> /etc/odbcinst.ini \
Expand All @@ -45,5 +44,4 @@ RUN echo '[ODBC Drivers]' > /etc/odbcinst.ini \
USER app

COPY src/api/ /home/site/wwwroot
COPY src /home/site/wwwroot/src

COPY src /home/site/wwwroot/src
4 changes: 1 addition & 3 deletions src/api/FastAPIApp/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,9 +65,7 @@
[ReDoc](/redoc)

[Real Time Data Ingestion Platform](https://www.rtdip.io/)
""".format(
os.environ.get("TENANT_ID")
)
""".format(os.environ.get("TENANT_ID"))

app = FastAPI(
title=TITLE,
Expand Down
2 changes: 1 addition & 1 deletion src/api/host.json
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,6 @@
},
"extensionBundle": {
"id": "Microsoft.Azure.Functions.ExtensionBundle",
"version": "[2.*, 3.0.0)"
"version": "[4.*, 5.0.0)"
}
}
3 changes: 2 additions & 1 deletion src/api/requirements.txt
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
# Do not include azure-functions-worker as it may conflict with the Azure Functions platform
setuptools>=68.0.0
azure-functions==1.20.0
fastapi==0.121.0
starlette==0.49.1
pydantic==2.10.1
# turbodbc==4.11.0
# turbodbc==4.5.3
pyodbc==5.2.0
importlib_metadata>=7.0.0,<9.0.0
databricks-sql-connector==3.6.0
Expand Down
3 changes: 1 addition & 2 deletions src/api/v1/batch.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,6 @@
from concurrent.futures import *
import pandas as pd


ROUTE_FUNCTION_MAPPING = {
"/events/raw": "raw",
"/events/latest": "latest",
Expand Down Expand Up @@ -120,7 +119,7 @@ async def batch_events_get(

try:
# Set up connection
(connection, parameters) = common_api_setup_tasks(
connection, _ = common_api_setup_tasks(
base_query_parameters=base_query_parameters,
base_headers=base_headers,
)
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/circular_average.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ def circular_average_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
raw_query_parameters=raw_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/circular_standard_deviation.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ def circular_standard_deviation_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
raw_query_parameters=raw_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
1 change: 0 additions & 1 deletion src/api/v1/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,6 @@

from src.sdk.python.rtdip_sdk.queries.time_series import batch


if importlib.util.find_spec("turbodbc") != None:
from src.sdk.python.rtdip_sdk.connectors import TURBODBCSQLConnection
from src.api.auth import azuread
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/interpolate.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ def interpolate_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
raw_query_parameters=raw_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/interpolation_at_time.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ def interpolation_at_time_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
tag_query_parameters=tag_query_parameters,
interpolation_at_time_query_parameters=interpolation_at_time_query_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/latest.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ def latest_retrieval_get(
query_parameters, metadata_query_parameters, limit_offset_parameters, base_headers
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
query_parameters,
metadata_query_parameters=metadata_query_parameters,
limit_offset_query_parameters=limit_offset_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ def metadata_retrieval_get(
query_parameters, metadata_query_parameters, limit_offset_parameters, base_headers
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
query_parameters,
metadata_query_parameters=metadata_query_parameters,
limit_offset_query_parameters=limit_offset_parameters,
Expand Down
11 changes: 5 additions & 6 deletions src/api/v1/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@
from src.api.auth.azuread import oauth2_scheme
from typing import Generic, TypeVar, Optional


EXAMPLE_DATE = "2022-01-01"
EXAMPLE_DATETIME = "2022-01-01T15:00:00"
EXAMPLE_DATETIME_TIMEZOME = "2022-01-01T15:00:00+00:00"
Expand Down Expand Up @@ -231,10 +230,10 @@ def __init__(
class BaseQueryParams:
def __init__(
self,
business_unit: str = Query(None, description="Business Unit Name"),
business_unit: str = Query(..., description="Business Unit Name"),
region: str = Query(..., description="Region"),
asset: str = Query(None, description="Asset"),
data_security_level: str = Query(None, description="Data Security Level"),
asset: str = Query(..., description="Asset"),
data_security_level: str = Query(..., description="Data Security Level"),
authorization: str = Depends(oauth2_scheme),
):
# Additional validation when mapping endpoint not provided - ensure validation error for missing params
Expand Down Expand Up @@ -300,7 +299,7 @@ class RawQueryParams:
def __init__(
self,
data_type: str = Query(
None,
...,
description="Data Type can be one of the following options: float, double, integer, string",
examples=["float", "double", "integer", "string"],
),
Expand Down Expand Up @@ -447,7 +446,7 @@ def __init__(
time_interval_rate: str = DuplicatedQueryParameters.time_interval_rate,
time_interval_unit: str = DuplicatedQueryParameters.time_interval_unit,
window_length: int = Query(
..., description="Window Length in days", examples=[1]
default=1, description="Window Length in days", examples=[1]
),
step: str = Query(
default="metadata",
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/plot.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ def plot_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
raw_query_parameters=raw_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/raw.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ def raw_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
raw_query_parameters=raw_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/resample.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ def resample_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
raw_query_parameters=raw_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/sql.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ def sql_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
sql_query_parameters=sql_query_parameters,
limit_offset_query_parameters=limit_offset_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/summary.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ def summary_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
summary_query_parameters=summary_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
2 changes: 1 addition & 1 deletion src/api/v1/time_weighted_average.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ def time_weighted_average_events_get(
base_headers,
):
try:
(connection, parameters) = common_api_setup_tasks(
connection, parameters = common_api_setup_tasks(
base_query_parameters,
raw_query_parameters=raw_query_parameters,
tag_query_parameters=tag_query_parameters,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,10 +49,8 @@ def validate_uri(uri: str):
return parsed_uri.scheme, parsed_uri.hostname, parsed_uri.path
except Exception as ex:
logging.error(ex)
raise SystemError(
f"Could not convert to valid tuple \
or scheme not supported: {uri} {parsed_uri.scheme}"
)
raise SystemError(f"Could not convert to valid tuple \
or scheme not supported: {uri} {parsed_uri.scheme}")


def get_supported_schema() -> list:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@
import hashlib
import time


series_id_str = "usage_series_id_001"
output_header_str: str = "Uid,SeriesId,Timestamp,IntervalTimestamp,Value"

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
import random
import logging


type_checks = [
# (Type, Test)
(int, int),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@
BooleanType,
)


RTDIP_FLOAT_WEATHER_DATA_MODEL = StructType(
[
StructField("TagName", StringType(), False),
Expand Down
4 changes: 2 additions & 2 deletions src/sdk/python/rtdip_sdk/pipelines/deploy/databricks.py
Original file line number Diff line number Diff line change
Expand Up @@ -395,7 +395,7 @@ def deploy(self) -> Union[bool, ValueError]:
module = self._load_module(
task.task_key + "file_upload", task.notebook_task.notebook_path
)
(task_libraries, spark_configuration) = PipelineComponentsGetUtility(
task_libraries, spark_configuration = PipelineComponentsGetUtility(
module.__name__
).execute()
workspace_client.workspace.mkdirs(path=self.workspace_directory)
Expand All @@ -415,7 +415,7 @@ def deploy(self) -> Union[bool, ValueError]:
module = self._load_module(
task.task_key + "file_upload", task.spark_python_task.python_file
)
(task_libraries, spark_configuration) = PipelineComponentsGetUtility(
task_libraries, spark_configuration = PipelineComponentsGetUtility(
module
).execute()
workspace_client.workspace.mkdirs(path=self.workspace_directory)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -105,15 +105,13 @@ def transform(self) -> DataFrame:
tag_name_expr.alias("TagName"),
lit("").alias("Description"),
lit("").alias("UoM"),
expr(
"""struct(
expr("""struct(
body.retroAltitude,
body.retroLongitude,
body.retroLatitude,
body.sensorAltitude,
body.sensorLongitude,
body.sensorLatitude)"""
).alias("Properties"),
body.sensorLatitude)""").alias("Properties"),
).dropDuplicates(["TagName"])

return df.select("TagName", "Description", "UoM", "Properties")
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ def settings() -> dict:
def execute(self) -> SparkSession:
"""To execute"""
try:
(task_libraries, spark_configuration) = PipelineComponentsGetUtility(
task_libraries, spark_configuration = PipelineComponentsGetUtility(
self.module, self.config
).execute()
self.spark = SparkClient(
Expand Down
Loading
Loading