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
4 changes: 2 additions & 2 deletions examples/rapidpro/contacts.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from api.rapidpro import pyRapid

contacts = pyRapid().contacts.get_contacts(
before="2023-01-02 00:00:00", after="2023-01-01 00:00:00"
end_datetime="2023-01-01 01:00:50", start_datetime="2023-01-01 00:00:00"
)

contacts.head(5)
print(contacts.collect())
2 changes: 1 addition & 1 deletion examples/rapidpro/fields.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,4 +2,4 @@

fields = pyRapid().fields.get_fields()

fields.head(5)
print(fields.collect())
2 changes: 1 addition & 1 deletion examples/rapidpro/flows.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,4 +2,4 @@

flows = pyRapid().flows.get_flows()

flows.head(5)
print(flows.collect())
4 changes: 2 additions & 2 deletions examples/rapidpro/flowstarts.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from api.rapidpro import pyRapid

flowstarts = pyRapid().flow_starts.get_flowstarts(
before="2023-01-02 00:00:00", after="2023-01-01 00:00:00"
end_datetime="2023-01-02 00:00:00", start_datetime="2023-01-01 00:00:00"
)

flowstarts.head(5)
print(flowstarts.collect())
2 changes: 1 addition & 1 deletion examples/rapidpro/groups.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,4 +2,4 @@

groups = pyRapid().groups.get_groups()

groups.head(5)
print(groups.collect())
4 changes: 2 additions & 2 deletions examples/rapidpro/runs.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from api.rapidpro import pyRapid

runs = pyRapid().runs.get_runs(
before="2024-10-01 01:30:00", after="2024-10-01 01:00:00"
end_datetime="2024-06-22 00:00:10", start_datetime="2024-06-22 00:00:00"
)

runs.head(5)
print(runs.collect())
6 changes: 6 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,15 @@ dev = [
"boto3-stubs>=1.37.19",
"starlette>=0.46.1",
"pytest-coverage>=0.0",
"pytest-mock>=3.14.1",
"werkzeug>=3.1.3",
]

[tool.coverage.run]
omit = [
"tests/*"
]

[build-system]
requires = ["hatchling"]
build-backend = "hatchling.build"
Expand Down
57 changes: 54 additions & 3 deletions rdw_ingestion_tools/api/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,15 @@
from collections.abc import Iterator

from pandas import DataFrame
from polars import LazyFrame, concat, json_normalize
from pandas import json_normalize as pd_json_normalize
from polars import (
LazyFrame,
Object,
String,
col,
concat,
json_normalize,
)

# https://github.com/astral-sh/ruff/issues/3388
from typing_extensions import Never # noqa: UP035
Expand Down Expand Up @@ -31,14 +39,14 @@ def concatenate(

"""
try:
df = concat([json_normalize(obj, sep="_") for obj in objs])
df = concat([pd_json_normalize(obj, sep="_") for obj in objs])
except ValueError:
df = DataFrame()

return df


def concatenate_to_lf(
def concatenate_to_lazyframe(
objs: list[dict] | dict[Never, Never] | list[Never] | Iterator, schema: dict
) -> LazyFrame:
"""
Expand All @@ -57,3 +65,46 @@ def concatenate_to_lf(
lf = LazyFrame(schema=schema)

return lf


def get_polars_schema(
object_columns: list[str], data: list[dict[str, Object]]
) -> dict[str, Object]:
"""
Creates a normalised LazyFrame and uses the schema to generate a schema
dictionary using the column names.
Columns that are `list` types need to be type `Object` before they can be cast
to string.
All other column types can be cast directly to string using the schema generated.
"""
# Create a dataframe to use to build the schema
columns = (
json_normalize(data, separator="_", infer_schema_length=None)
.lazy()
.collect_schema()
.names()
)
schema = {
column: (Object if column in object_columns else String) for column in columns
}

return schema


def concatenate_to_string_lazyframe(
objs: list[dict] | dict[Never, Never] | list[Never] | Iterator,
object_columns: list[str],
) -> LazyFrame:
"""
Flattens JSON data. Returns a LazyFrame with String columns.
"""
data = list(objs)

schema = get_polars_schema(data=data, object_columns=object_columns)
lf = (
json_normalize(data, separator="_", schema=schema)
.lazy()
.with_columns(col(Object).map_elements(lambda x: str(x), return_dtype=String))
)

return lf
8 changes: 6 additions & 2 deletions rdw_ingestion_tools/api/rapidpro/extensions/httpx.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
from collections.abc import Iterator
from urllib.parse import unquote

from httpx import Client

Expand All @@ -14,9 +15,10 @@ def get_paginated(
pages are left to return.

"""
params = {**kwargs}

while True:
response = client.get(url, params={**kwargs})
response = client.get(url, params=params)
response.raise_for_status()

data: dict = response.json()
Expand All @@ -25,6 +27,8 @@ def get_paginated(
yield from results

try:
url = data["next"].split("/v2/")[1]
cursor = data["next"].split("cursor=")[1].split("&")[0]
decoded_cursor = unquote(cursor)
params["cursor"] = decoded_cursor
except AttributeError:
break
22 changes: 14 additions & 8 deletions rdw_ingestion_tools/api/rapidpro/requests/contacts.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
from attrs import define
from httpx import Client
from pandas import DataFrame
from polars import LazyFrame

from api import concatenate
from api import concatenate_to_string_lazyframe

from ..extensions.httpx import get_paginated

Expand All @@ -13,22 +13,28 @@ class Contacts:

client: Client

def get_contacts(self, **kwargs: str | int) -> DataFrame:
"""Get a pandas DataFrame of Rapidpro contacts.
def get_contacts(
self, start_datetime: str, end_datetime: str, **kwargs: str | int
) -> LazyFrame:
"""Get a Polars LazyFrame of Rapidpro contacts.

This endpoint supports time-based filtering that allows
to fetch results between two date parameters. Example:

pyRapid.contacts.get_contacts(
before="2023-01-02T00:00:00",
after="2023-01-01T00:00:00"
end_datetime="2023-01-02T00:00:00",
start_datetime="2023-01-01T00:00:00"
)

"""
url = "contacts.json"

contacts_generator = get_paginated(self.client, url, **kwargs)
contacts_generator = get_paginated(
self.client, url, after=start_datetime, before=end_datetime, **kwargs
)

contacts = concatenate(contacts_generator)
contacts = concatenate_to_string_lazyframe(
contacts_generator, object_columns=["urns", "groups"]
)

return contacts
10 changes: 5 additions & 5 deletions rdw_ingestion_tools/api/rapidpro/requests/fields.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
from attrs import define
from httpx import Client
from pandas import DataFrame
from polars import LazyFrame

from api import concatenate
from api import concatenate_to_string_lazyframe

from ..extensions.httpx import get_paginated

Expand All @@ -13,8 +13,8 @@ class Fields:

client: Client

def get_fields(self, **kwargs: str | int) -> DataFrame:
"""Get a pandas DataFrame of Rapidpro fields.
def get_fields(self, **kwargs: str | int) -> LazyFrame:
"""Get a polars LazyFrame of Rapidpro fields.

This endpoint does not support time-based filtering and
can be called as:
Expand All @@ -26,6 +26,6 @@ def get_fields(self, **kwargs: str | int) -> DataFrame:

fields_generator = get_paginated(self.client, url, **kwargs)

fields = concatenate(fields_generator)
fields = concatenate_to_string_lazyframe(fields_generator, object_columns=[])

return fields
22 changes: 14 additions & 8 deletions rdw_ingestion_tools/api/rapidpro/requests/flow_starts.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
from attrs import define
from httpx import Client
from pandas import DataFrame
from polars import LazyFrame

from api import concatenate
from api import concatenate_to_string_lazyframe

from ..extensions.httpx import get_paginated

Expand All @@ -13,22 +13,28 @@ class FlowStarts:

client: Client

def get_flowstarts(self, **kwargs: str | int) -> DataFrame:
"""Get a pandas DataFrame of Rapidpro flowstarts.
def get_flowstarts(
self, start_datetime: str, end_datetime: str, **kwargs: str | int
) -> LazyFrame:
"""Get a Polars LazyFrame of Rapidpro flowstarts.

This endpoint supports time-based filtering that allows
to fetch results between two date parameters. Example:

pyRapid.flowstarts.get_flowstarts(
before="2023-01-02T00:00:00",
after="2023-01-01T00:00:00"
end_datetime="2023-01-02T00:00:00",
start_datetime="2023-01-01T00:00:00"
)

"""
url = "flow_starts.json"

flowstarts_generator = get_paginated(self.client, url, **kwargs)
flowstarts_generator = get_paginated(
self.client, url, after=start_datetime, before=end_datetime, **kwargs
)

flowstarts = concatenate(flowstarts_generator)
flowstarts = concatenate_to_string_lazyframe(
flowstarts_generator, object_columns=["groups", "contacts"]
)

return flowstarts
12 changes: 7 additions & 5 deletions rdw_ingestion_tools/api/rapidpro/requests/flows.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
from attrs import define
from httpx import Client
from pandas import DataFrame
from polars import LazyFrame

from api import concatenate
from api import concatenate_to_string_lazyframe

from ..extensions.httpx import get_paginated

Expand All @@ -13,8 +13,8 @@ class Flows:

client: Client

def get_flows(self, **kwargs: str | int) -> DataFrame:
"""Get a pandas DataFrame of Rapidpro flows.
def get_flows(self, **kwargs: str | int) -> LazyFrame:
"""Get a Polars LazyFrame of Rapidpro flows.

This endpoint does not support time-based filtering and
can be called as:
Expand All @@ -26,6 +26,8 @@ def get_flows(self, **kwargs: str | int) -> DataFrame:

flows_generator = get_paginated(self.client, url, **kwargs)

flows = concatenate(flows_generator)
flows = concatenate_to_string_lazyframe(
flows_generator, object_columns=["labels", "results", "parent_refs"]
)

return flows
10 changes: 5 additions & 5 deletions rdw_ingestion_tools/api/rapidpro/requests/groups.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
from attrs import define
from httpx import Client
from pandas import DataFrame
from polars import LazyFrame

from api import concatenate
from api import concatenate_to_string_lazyframe

from ..extensions.httpx import get_paginated

Expand All @@ -13,8 +13,8 @@ class Groups:

client: Client

def get_groups(self, **kwargs: str | int) -> DataFrame:
"""Get a pandas DataFrame of Rapidpro groups.
def get_groups(self, **kwargs: str | int) -> LazyFrame:
"""Get a Polars LazyFrame of Rapidpro groups.

This endpoint does not support time-based filtering and
can be called as:
Expand All @@ -26,6 +26,6 @@ def get_groups(self, **kwargs: str | int) -> DataFrame:

groups_generator = get_paginated(self.client, url, **kwargs)

groups = concatenate(groups_generator)
groups = concatenate_to_string_lazyframe(groups_generator, object_columns=[])

return groups
Loading