Skip to content
Merged

Dev #11

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 coded-flows/coded_flows/types/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -488,8 +488,8 @@ def is_supported_type(element_type):
"Any",
"Null",
# Data
"DataSeries", # <-- works as a Helper
"DataFrame", # <-- works as a Helper
"DataSeries",
"DataFrame",
"ArrowTable",
"NDArray",
"DataDict",
Expand Down
78 changes: 63 additions & 15 deletions coded-flows/coded_flows/types/extra.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,38 +2,86 @@
import base64
import pandas as pd
import pyarrow as pa
import polars as pl
from numpy import ndarray
from pydantic import GetCoreSchemaHandler
from pydantic_core import core_schema
from typing import Any, Type
from typing import Any, Type, Union, List
from PIL import Image


class DataSeries(pd.Series):
class DataSeriesMeta(type):

def __instancecheck__(cls, instance):
return isinstance(instance, (pd.Series, pl.Series))


def serialize_series(series: Union[pd.Series, pl.Series]) -> List[Any]:
return series.to_list()


class DataSeries(metaclass=DataSeriesMeta):
@classmethod
def __get_pydantic_core_schema__(
cls, _source: Type[Any], _handler: GetCoreSchemaHandler
cls,
source: Type[Any],
handler: GetCoreSchemaHandler,
) -> core_schema.CoreSchema:
pandas_schema = core_schema.is_instance_schema(pd.Series)
polars_schema = core_schema.is_instance_schema(pl.Series)

return core_schema.is_instance_schema(
pd.Series,
serialization=core_schema.plain_serializer_function_ser_schema(
lambda instance: list(instance)
),
union_schema = core_schema.union_schema([pandas_schema, polars_schema])

serialization = core_schema.plain_serializer_function_ser_schema(
serialize_series, when_used="json"
)

return core_schema.json_or_python_schema(
json_schema=union_schema,
python_schema=union_schema,
serialization=serialization,
)


class DataFrame(pd.DataFrame):
class DataFrameMeta(type):
def __instancecheck__(cls, instance):
return isinstance(instance, (pd.DataFrame, pl.DataFrame, pl.LazyFrame))


def serialize_dataframe(
df: Union[pd.DataFrame, pl.DataFrame, pl.LazyFrame],
) -> list[dict[str, Any]]:
if isinstance(df, pd.DataFrame):
return df.to_dict(orient="records")
if isinstance(df, pl.DataFrame):
return df.to_dicts()
if isinstance(df, pl.LazyFrame):
return df.collect().to_dicts()
raise TypeError(f"Unsupported dataframe type: {type(df)}")


class DataFrame(metaclass=DataFrameMeta):
@classmethod
def __get_pydantic_core_schema__(
cls, _source: Type[Any], _handler: GetCoreSchemaHandler
cls,
_source: Type[Any],
_handler: GetCoreSchemaHandler,
) -> core_schema.CoreSchema:
pandas_schema = core_schema.is_instance_schema(pd.DataFrame)
polars_eager_schema = core_schema.is_instance_schema(pl.DataFrame)
polars_lazy_schema = core_schema.is_instance_schema(pl.LazyFrame)
union_schema = core_schema.union_schema(
[pandas_schema, polars_eager_schema, polars_lazy_schema]
)

return core_schema.is_instance_schema(
pd.DataFrame,
serialization=core_schema.plain_serializer_function_ser_schema(
lambda instance: instance.to_dict(orient="records")
),
serialization = core_schema.plain_serializer_function_ser_schema(
serialize_dataframe, when_used="json"
)

return core_schema.json_or_python_schema(
json_schema=union_schema,
python_schema=union_schema,
serialization=serialization,
)


Expand Down
56 changes: 47 additions & 9 deletions coded-flows/coded_flows/utils/converters.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from typing import Any, Callable, Union
from pydantic_core import MultiHostUrl
import pandas as pd
import polars as pl
import pyarrow as pa
from ..types import (
AnyUrl,
Expand Down Expand Up @@ -260,28 +261,65 @@ def dataseries_to_type(output_type: str) -> Callable:
elif output_type == "Set":
return lambda x: set(x.to_list())
elif output_type == "Json":
return lambda x: x.to_json()
return lambda x: (
json.dumps(x.to_list())
if isinstance(x, pl.Series)
else x.to_json(orient="records")
)
elif output_type == "NDArray":
return lambda x: x.to_numpy()
elif output_type == "ArrowTable":
return lambda x: pa.Table.from_pandas(x.to_frame())
return lambda x: (
pa.Table.from_arrays([x.to_arrow()], names=[x.name if x.name else "value"])
if isinstance(x, pl.Series)
else pa.Table.from_pandas(x.to_frame())
)


def _collect_if_lazy(df: pl.DataFrame | pl.LazyFrame | pd.DataFrame) -> pl.DataFrame:
return df.collect() if isinstance(df, pl.LazyFrame) else df


def dataframe_to_type(output_type: str) -> Callable:

if output_type == "DataRecords":
return lambda x: x.to_dict("records")
return lambda x: (
_collect_if_lazy(x).to_dicts()
if isinstance(_collect_if_lazy(x), pl.DataFrame)
else x.to_dict("records")
)
elif output_type == "List":
return lambda x: x.to_dict("records")
return lambda x: (
_collect_if_lazy(x).to_dicts()
if isinstance(_collect_if_lazy(x), pl.DataFrame)
else x.to_dict("records")
)
elif output_type == "Dict":
return lambda x: x.to_dict("list")
return lambda x: (
_collect_if_lazy(x).to_dict(as_series=False)
if isinstance(_collect_if_lazy(x), pl.DataFrame)
else x.to_dict("list")
)
elif output_type == "DataDict":
return lambda x: x.to_dict("list")
return lambda x: (
_collect_if_lazy(x).to_dict(as_series=False)
if isinstance(_collect_if_lazy(x), pl.DataFrame)
else x.to_dict("list")
)
elif output_type == "Json":
return lambda x: x.to_json(orient="records")
return lambda x: (
_collect_if_lazy(x).write_json()
if isinstance(_collect_if_lazy(x), pl.DataFrame)
else x.to_json(orient="records")
)
elif output_type == "NDArray":
return lambda x: x.to_numpy()
return lambda x: _collect_if_lazy(x).to_numpy()
elif output_type == "ArrowTable":
return lambda x: pa.Table.from_pandas(x)
return lambda x: (
_collect_if_lazy(x).to_arrow()
if isinstance(_collect_if_lazy(x), pl.DataFrame)
else pa.Table.from_pandas(x)
)


def arrow_to_type(output_type: str) -> Callable:
Expand Down
34 changes: 25 additions & 9 deletions coded-flows/coded_flows/utils/media.py
Original file line number Diff line number Diff line change
Expand Up @@ -147,27 +147,29 @@ def save_data_to_parquet(
filename=None,
) -> str:

random_filename = f"cfdata_{filename if filename else uuid.uuid4().hex}.parquet"
temp_dir = os.path.join(tempfile.gettempdir(), "coded-flows-media")
os.makedirs(temp_dir, exist_ok=True)
file_path = os.path.join(temp_dir, random_filename)

if (
isinstance(data, pd.DataFrame)
or isinstance(data, pd.Series)
or isinstance(data, pa.Table)
or (
isinstance(data, list) and all(isinstance(item, dict) for item in data[:50])
)
):

random_filename = f"cfdata_{filename if filename else uuid.uuid4().hex}.parquet"
temp_dir = os.path.join(tempfile.gettempdir(), "coded-flows-media")
os.makedirs(temp_dir, exist_ok=True)
file_path = os.path.join(temp_dir, random_filename)

try:

# Handle different data types
if isinstance(data, pd.DataFrame):
data.to_parquet(
file_path, row_group_size=50000, index=False, engine="pyarrow"
)

elif isinstance(data, pd.Series):
data.to_frame().to_parquet(
file_path, row_group_size=50000, index=False, engine="pyarrow"
)
elif isinstance(data, pa.Table):
pq.write_table(data, file_path, row_group_size=50000)

Expand All @@ -180,4 +182,18 @@ def save_data_to_parquet(
except Exception as e:
raise Exception(f"❌ Error saving data to parquet: {str(e)}")
else:
raise TypeError(f"Unsupported data type: {type(data)}")

if isinstance(data, np.ndarray):
if data.ndim != 1:
raise ValueError(f"NumPy array must be one-dimensional.")
pd.DataFrame(data, columns=["value"]).to_parquet(
file_path, row_group_size=50000, index=False, engine="pyarrow"
)
elif isinstance(data, list):
pd.DataFrame(data, columns=["value"]).to_parquet(
file_path, row_group_size=50000, index=False, engine="pyarrow"
)
else:
raise TypeError(f"Unsupported data type: {type(data)}")

return file_path
46 changes: 45 additions & 1 deletion coded-flows/poetry.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 2 additions & 1 deletion coded-flows/pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[tool.poetry]
name = "coded-flows"
version = "0.6.1"
version = "0.6.2"
description = "Various utilities for Coded Flows"
authors = ["COLOR CODED CODES <contact@colorcoded.codes>"]
readme = "README.md"
Expand All @@ -14,6 +14,7 @@ email-validator = "^2.0.0"
pycountry = "^24.6.1"
pillow = "^10.1.0"
pydantic-extra-types = "^2.6.0"
polars = "^1.0.0"

[tool.poetry.group.dev.dependencies]
pytest = "^8.1.1"
Expand Down
Loading