Skip to content

Commit f29ef62

Browse files
committed
Remove the need for a mapper/sinker superclass. Just use callables
Signed-off-by: Sreekanth <prsreekanth920@gmail.com>
1 parent 9f664b2 commit f29ef62

22 files changed

Lines changed: 82 additions & 167 deletions

packages/pynumaflow-lite/manifests/batchmap/batchmap_cat.py

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,10 @@
11
from collections.abc import AsyncIterator
22

33
from pynumaflow_lite import batchmapper
4-
from pynumaflow_lite.batchmapper import BatchMapper, BatchResponse, Datum, Message
4+
from pynumaflow_lite.batchmapper import BatchResponse, Datum, Message
55

66

7-
class SimpleBatchCat(BatchMapper):
7+
class SimpleBatchCat:
88
async def handler(self, batch: AsyncIterator[Datum]) -> list[BatchResponse]:
99
return [
1010
BatchResponse(
@@ -16,4 +16,5 @@ async def handler(self, batch: AsyncIterator[Datum]) -> list[BatchResponse]:
1616

1717

1818
if __name__ == "__main__":
19-
batchmapper.BatchMapAsyncServer(SimpleBatchCat()).run()
19+
batch_mapper_obj = SimpleBatchCat()
20+
batchmapper.BatchMapAsyncServer(batch_mapper_obj.handler).run()
Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,13 @@
11
from pynumaflow_lite import mapper
22

33

4-
class SimpleCat(mapper.Mapper):
4+
class SimpleCat:
55
async def handler(self, datum: mapper.Datum) -> list[mapper.Message]:
66
if datum.value == b"bad world":
77
return [mapper.Message.to_drop()]
88
return [mapper.Message(datum.value, keys=datum.keys)]
99

1010

1111
if __name__ == "__main__":
12-
mapper.MapAsyncServer(SimpleCat()).run()
12+
mapper_obj = SimpleCat()
13+
mapper.MapAsyncServer(mapper_obj.handler).run()

packages/pynumaflow-lite/manifests/mapstream/mapstream_cat.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44
from pynumaflow_lite.mapstreamer import Message
55

66

7-
class SimpleStreamCat(mapstreamer.MapStreamer):
7+
class SimpleStreamCat:
88
async def handler(self, datum: mapstreamer.Datum) -> AsyncIterable[Message]:
99
parts = datum.value.decode("utf-8").split(",")
1010
if not parts:
@@ -15,4 +15,5 @@ async def handler(self, datum: mapstreamer.Datum) -> AsyncIterable[Message]:
1515

1616

1717
if __name__ == "__main__":
18-
mapstreamer.MapStreamAsyncServer(SimpleStreamCat()).run()
18+
map_streamer_obj = SimpleStreamCat()
19+
mapstreamer.MapStreamAsyncServer(map_streamer_obj.handler).run()

packages/pynumaflow-lite/manifests/sideinput/sideinput_example.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ async def retrieve_handler(self) -> sideinputer.Response:
3939
return sideinputer.Response.broadcast_message(val.encode("utf-8"))
4040

4141

42-
class SideInputHandler(mapper.Mapper):
42+
class SideInputHandler:
4343
"""
4444
A Mapper that reads from side input files and includes the value in its output.
4545
"""
@@ -103,16 +103,16 @@ async def start_sideinput():
103103

104104
def start_mapper():
105105
"""Start the Mapper server that reads from side inputs."""
106-
handler = SideInputHandler()
106+
mapper_obj = SideInputHandler()
107107

108108
# Initialize the data value from the side input file
109-
handler.init_data_value()
109+
mapper_obj.init_data_value()
110110

111111
# Start the file watcher in a background thread
112-
watcher_thread = Thread(target=handler.file_watcher, daemon=True)
112+
watcher_thread = Thread(target=mapper_obj.file_watcher, daemon=True)
113113
watcher_thread.start()
114114

115-
mapper.MapAsyncServer(handler).run()
115+
mapper.MapAsyncServer(mapper_obj.handler).run()
116116

117117

118118
if __name__ == "__main__":
Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,19 @@
11
import logging
2-
from collections.abc import AsyncIterable
2+
from collections.abc import AsyncIterator
33

44
from pynumaflow_lite import sinker
5-
from pynumaflow_lite.sinker import Sinker
65

76
# Configure logging
87
logging.basicConfig(level=logging.INFO)
98
_LOGGER = logging.getLogger(__name__)
109

1110

12-
class SimpleLogSink(Sinker):
11+
class SimpleLogSink:
1312
"""
1413
Simple log sink that logs each message and returns success responses.
1514
"""
1615

17-
async def handler(self, datums: AsyncIterable[sinker.Datum]) -> list[sinker.Response]:
16+
async def handler(self, datums: AsyncIterator[sinker.Datum]) -> list[sinker.Response]:
1817
responses = []
1918
async for msg in datums:
2019
_LOGGER.info("User Defined Sink: %s", msg.value.decode("utf-8"))
@@ -25,4 +24,5 @@ async def handler(self, datums: AsyncIterable[sinker.Datum]) -> list[sinker.Resp
2524

2625

2726
if __name__ == "__main__":
28-
sinker.SinkAsyncServer(SimpleLogSink()).run()
27+
sinker_obj = SimpleLogSink()
28+
sinker.SinkAsyncServer(sinker_obj.handler).run()

packages/pynumaflow-lite/pynumaflow_lite/__init__.py

Lines changed: 1 addition & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -61,34 +61,26 @@
6161
except Exception: # pragma: no cover
6262
sideinputer = None
6363

64-
# Surface the Python Mapper, BatchMapper, MapStreamer, Reducer, SessionReducer, ReduceStreamer, Accumulator, Sinker,
65-
# Sourcer, SourceTransformer, and SideInput classes under the extension submodules for convenient access
64+
# Surface the Python async servers and data-type classes under the extension submodules for convenient access
6665
from ._accumulator_dtypes import Accumulator
6766
from ._batchmap_server import BatchMapAsyncServer
68-
from ._batchmapper_dtypes import BatchMapper
69-
from ._map_dtypes import Mapper
7067
from ._map_server import MapAsyncServer
71-
from ._mapstream_dtypes import MapStreamer
7268
from ._mapstream_server import MapStreamAsyncServer
7369
from ._reduce_dtypes import Reducer
7470
from ._reducestreamer_dtypes import ReduceStreamer
7571
from ._session_reduce_dtypes import SessionReducer
7672
from ._sideinput_dtypes import SideInput
77-
from ._sink_dtypes import Sinker
7873
from ._sink_server import SinkAsyncServer
7974
from ._source_dtypes import Sourcer
8075
from ._sourcetransformer_dtypes import SourceTransformer
8176

8277
if mapper is not None:
83-
mapper.Mapper = Mapper
8478
mapper.MapAsyncServer = MapAsyncServer
8579

8680
if batchmapper is not None:
87-
batchmapper.BatchMapper = BatchMapper
8881
batchmapper.BatchMapAsyncServer = BatchMapAsyncServer
8982

9083
if mapstreamer is not None:
91-
mapstreamer.MapStreamer = MapStreamer
9284
mapstreamer.MapStreamAsyncServer = MapStreamAsyncServer
9385

9486
if reducer is not None:
@@ -104,7 +96,6 @@
10496
accumulator.Accumulator = Accumulator
10597

10698
if sinker is not None:
107-
sinker.Sinker = Sinker
10899
sinker.SinkAsyncServer = SinkAsyncServer
109100

110101
if sourcer is not None:

packages/pynumaflow-lite/pynumaflow_lite/_batchmap_server.py

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,16 +2,22 @@
22

33
import asyncio
44
import signal
5+
from collections.abc import AsyncIterator, Awaitable, Callable
56
from types import TracebackType
6-
from typing import Any
77

88
from .pynumaflow_lite import batchmapper as _batchmapper
99

10+
BatchResponse = _batchmapper.BatchResponse
11+
Datum = _batchmapper.Datum
12+
1013

1114
class BatchMapAsyncServer:
1215
def __init__(
1316
self,
14-
handler: Any,
17+
handler: Callable[
18+
[AsyncIterator[Datum]],
19+
Awaitable[list[BatchResponse]],
20+
],
1521
*,
1622
sock_file: str | None = None,
1723
server_info_file: str | None = None,

packages/pynumaflow-lite/pynumaflow_lite/_batchmapper_dtypes.py

Lines changed: 0 additions & 21 deletions
This file was deleted.

packages/pynumaflow-lite/pynumaflow_lite/_map_dtypes.py

Lines changed: 0 additions & 20 deletions
This file was deleted.

packages/pynumaflow-lite/pynumaflow_lite/_map_server.py

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,16 +2,19 @@
22

33
import asyncio
44
import signal
5+
from collections.abc import Awaitable, Callable
56
from types import TracebackType
6-
from typing import Any
77

88
from .pynumaflow_lite import mapper as _mapper
99

10+
Datum = _mapper.Datum
11+
Message = _mapper.Message
12+
1013

1114
class MapAsyncServer:
1215
def __init__(
1316
self,
14-
handler: Any,
17+
handler: Callable[[Datum], Awaitable[list[Message]]],
1518
*,
1619
sock_file: str | None = None,
1720
server_info_file: str | None = None,

0 commit comments

Comments
 (0)