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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
.venv/
80 changes: 43 additions & 37 deletions packages/pynumaflow-lite/manifests/batchmap/Dockerfile
Original file line number Diff line number Diff line change
@@ -1,37 +1,43 @@
FROM python:3.11-slim-bullseye AS builder

ENV PYTHONFAULTHANDLER=1 \
PYTHONUNBUFFERED=1 \
PYTHONHASHSEED=random \
PIP_NO_CACHE_DIR=on \
PIP_DISABLE_PIP_VERSION_CHECK=on \
PIP_DEFAULT_TIMEOUT=100 \
POETRY_HOME="/opt/poetry" \
POETRY_VIRTUALENVS_IN_PROJECT=true \
POETRY_NO_INTERACTION=1 \
PYSETUP_PATH="/opt/pysetup"

ENV PATH="$POETRY_HOME/bin:$PATH"

RUN apt-get update \
&& apt-get install --no-install-recommends -y \
curl \
wget \
# deps for building python deps
build-essential \
&& apt-get install -y git \
&& apt-get clean && rm -rf /var/lib/apt/lists/* \
&& curl -sSL https://install.python-poetry.org | python3 -

FROM builder AS udf

WORKDIR $PYSETUP_PATH
COPY ./ ./

RUN pip install $PYSETUP_PATH/pynumaflow_lite-0.1.0-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl

RUN poetry lock
RUN poetry install --no-cache --no-root && \
rm -rf ~/.cache/pypoetry/

CMD ["python", "batchmap_cat.py"]
FROM python:3.11-slim-trixie AS builder

COPY --from=ghcr.io/astral-sh/uv:0.12.13 /uv /uvx /bin/

ENV UV_COMPILE_BYTECODE=1 \
UV_LINK_MODE=copy \
UV_PYTHON_DOWNLOADS=never

WORKDIR /app

RUN --mount=type=cache,target=/root/.cache/uv \
--mount=type=bind,source=uv.lock,target=uv.lock \
--mount=type=bind,source=pyproject.toml,target=pyproject.toml \
uv sync --locked --no-install-project --no-dev

COPY . .

RUN --mount=type=cache,target=/root/.cache/uv \
uv sync --locked --no-dev

# Install the local pynumaflow-lite wheel that matches the build platform.
RUN uv pip install --no-index --find-links . pynumaflow-lite

FROM python:3.11-slim-trixie

# Setup a non-root user
RUN groupadd --system --gid 999 nonroot \
&& useradd --system --gid 999 --uid 999 --create-home nonroot

COPY --from=builder --chown=nonroot:nonroot /app /app

ENV PATH="/app/.venv/bin:$PATH"

# Keeps Python from buffering stdout and stderr to avoid situations where
# the application crashes without emitting any logs due to buffering.
ENV PYTHONUNBUFFERED=1

# Use the non-root user to run our application
USER nonroot

WORKDIR /app

CMD ["python", "batchmap_cat.py"]
52 changes: 16 additions & 36 deletions packages/pynumaflow-lite/manifests/batchmap/batchmap_cat.py
Original file line number Diff line number Diff line change
@@ -1,46 +1,26 @@
import asyncio
import signal
from collections.abc import AsyncIterable, Awaitable, Callable
from collections.abc import AsyncIterable

from pynumaflow_lite import batchmapper
from pynumaflow_lite.batchmapper import Message
from pynumaflow_lite.batchmapper import BatchMapAsyncServer, BatchMapper, BatchResponse, Datum, Message


class SimpleBatchCat(batchmapper.BatchMapper):
async def handler(self, batch: AsyncIterable[batchmapper.Datum]) -> batchmapper.BatchResponses:
responses = batchmapper.BatchResponses()
async for d in batch:
resp = batchmapper.BatchResponse(d.id)
if d.value == b"bad world":
resp.append(Message.message_to_drop())
continue

resp.append(Message(d.value, d.keys))
responses.append(resp)
class SimpleBatchCat(BatchMapper):
async def handler(self, batch: AsyncIterable[Datum]) -> list[BatchResponse]:
responses = []
async for datum in batch:
if datum.value == b"bad world":
responses.append(BatchResponse(datum.id, Message.to_drop()))
else:
responses.append(BatchResponse(datum.id, Message(datum.value, keys=datum.keys)))
return responses


async def start(
f: Callable[[AsyncIterable[batchmapper.Datum]], Awaitable[batchmapper.BatchResponses]],
):
server = batchmapper.BatchMapAsyncServer()

# Register loop-level signal handlers so we control shutdown and avoid asyncio.run
loop = asyncio.get_running_loop()
try:
loop.add_signal_handler(signal.SIGINT, lambda: server.stop())
loop.add_signal_handler(signal.SIGTERM, lambda: server.stop())
except (NotImplementedError, RuntimeError):
pass

try:
await server.start(f)
print("Shutting down gracefully...")
except asyncio.CancelledError:
server.stop()
return
async def main() -> None:
print("Starting BatchMap server")
# `serve` returns when SIGINT or SIGTERM arrives.
await BatchMapAsyncServer(SimpleBatchCat()).serve()
print("BatchMap server stopped")


if __name__ == "__main__":
async_handler = SimpleBatchCat()
asyncio.run(start(async_handler))
asyncio.run(main())
6 changes: 3 additions & 3 deletions packages/pynumaflow-lite/manifests/batchmap/pipeline.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@ spec:
source:
# A self data generating source
generator:
rpu: 500
duration: 1s
rpu: 3
duration: 3s
- name: batchmap
partitions: 2
scale:
Expand All @@ -27,4 +27,4 @@ spec:
- from: in
to: batchmap
- from: batchmap
to: sink
to: sink
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,6 @@ authors = [
{ name = "Vigith Maurice", email = "vigith@gmail.com" }
]
readme = "README.md"
requires-python = ">=3.11"
requires-python = "==3.11.*"
dependencies = [
]


[build-system]
requires = ["poetry-core>=2.0.0,<3.0.0"]
build-backend = "poetry.core.masonry.api"
8 changes: 8 additions & 0 deletions packages/pynumaflow-lite/manifests/batchmap/uv.lock

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

126 changes: 31 additions & 95 deletions packages/pynumaflow-lite/pynumaflow_lite/__init__.py
Original file line number Diff line number Diff line change
@@ -1,69 +1,11 @@
from . import (
pynumaflow_lite, # type: ignore[attr-defined] # Rust extension, resolved at runtime
)
from .pynumaflow_lite import * # noqa: F403 # Rust extension; exports resolved at runtime

# Ensure the `mapper`, `batchmapper`, and `mapstreamer` submodules are importable as attributes of the package
# even though they're primarily registered by the extension module.
try:
from importlib import import_module as _import_module

mapper = _import_module(__name__ + ".mapper")
except Exception: # pragma: no cover - avoid hard failures if extension not built
mapper = None

try:
batchmapper = _import_module(__name__ + ".batchmapper")
except Exception: # pragma: no cover
batchmapper = None

try:
mapstreamer = _import_module(__name__ + ".mapstreamer")
except Exception: # pragma: no cover
mapstreamer = None
try:
reducer = _import_module(__name__ + ".reducer")
except Exception: # pragma: no cover
reducer = None

try:
session_reducer = _import_module(__name__ + ".session_reducer")
except Exception: # pragma: no cover
session_reducer = None

try:
reducestreamer = _import_module(__name__ + ".reducestreamer")
except Exception: # pragma: no cover
reducestreamer = None

try:
accumulator = _import_module(__name__ + ".accumulator")
except Exception: # pragma: no cover
accumulator = None

try:
sinker = _import_module(__name__ + ".sinker")
except Exception: # pragma: no cover
sinker = None

try:
sourcer = _import_module(__name__ + ".sourcer")
except Exception: # pragma: no cover
sourcer = None

try:
sourcetransformer = _import_module(__name__ + ".sourcetransformer")
except Exception: # pragma: no cover
sourcetransformer = None

try:
sideinputer = _import_module(__name__ + ".sideinputer")
except Exception: # pragma: no cover
sideinputer = None

# Surface the Python Mapper, BatchMapper, MapStreamer, Reducer, SessionReducer, ReduceStreamer, Accumulator, Sinker,
# Sourcer, SourceTransformer, and SideInput classes under the extension submodules for convenient access
from ._accumulator_dtypes import Accumulator
from ._batchmap_server import BatchMapAsyncServer
from ._batchmapper_dtypes import BatchMapper
from ._map_dtypes import Mapper
from ._map_server import MapAsyncServer
Expand All @@ -76,41 +18,38 @@
from ._sink_server import SinkAsyncServer
from ._source_dtypes import Sourcer
from ._sourcetransformer_dtypes import SourceTransformer
from .pynumaflow_lite import * # noqa: F403 # Rust extension; exports resolved at runtime

if mapper is not None:
mapper.Mapper = Mapper
mapper.MapAsyncServer = MapAsyncServer

if batchmapper is not None:
batchmapper.BatchMapper = BatchMapper

if mapstreamer is not None:
mapstreamer.MapStreamer = MapStreamer

if reducer is not None:
reducer.Reducer = Reducer

if session_reducer is not None:
session_reducer.SessionReducer = SessionReducer

if reducestreamer is not None:
reducestreamer.ReduceStreamer = ReduceStreamer

if accumulator is not None:
accumulator.Accumulator = Accumulator

if sinker is not None:
sinker.Sinker = Sinker
sinker.SinkAsyncServer = SinkAsyncServer

if sourcer is not None:
sourcer.Sourcer = Sourcer

if sourcetransformer is not None:
sourcetransformer.SourceTransformer = SourceTransformer
# Submodules are defined by the Rust extension, which also registers them in sys.modules
# as `pynumaflow_lite.<name>`.
from .pynumaflow_lite import ( # type: ignore[attr-defined]
accumulator,
batchmapper,
mapper,
mapstreamer,
reducer,
reducestreamer,
session_reducer,
sideinputer,
sinker,
sourcer,
sourcetransformer,
)

if sideinputer is not None:
sideinputer.SideInput = SideInput
mapper.Mapper = Mapper
mapper.MapAsyncServer = MapAsyncServer
batchmapper.BatchMapper = BatchMapper
batchmapper.BatchMapAsyncServer = BatchMapAsyncServer
mapstreamer.MapStreamer = MapStreamer
reducer.Reducer = Reducer
session_reducer.SessionReducer = SessionReducer
reducestreamer.ReduceStreamer = ReduceStreamer
accumulator.Accumulator = Accumulator
sinker.Sinker = Sinker
sinker.SinkAsyncServer = SinkAsyncServer
sourcer.Sourcer = Sourcer
sourcetransformer.SourceTransformer = SourceTransformer
sideinputer.SideInput = SideInput

# Public API
__all__ = [
Expand All @@ -128,6 +67,3 @@
]

__doc__ = pynumaflow_lite.__doc__
if hasattr(pynumaflow_lite, "__all__"):

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

all names are already added in as a static list above

# Merge to keep our package-level exports
__all__ = list(set(__all__) | set(pynumaflow_lite.__all__))
Loading
Loading