Skip to content

Commit 1079b70

Browse files
feat: introduce API to fail a message in udf/transformer (#370)
Signed-off-by: Vaibhav Tiwari <vaibhav.tiwari33@gmail.com>
1 parent e96694e commit 1079b70

17 files changed

Lines changed: 176 additions & 14 deletions

File tree

packages/pynumaflow/pynumaflow/_constants.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@
5555
DELIMITER = ":"
5656
DROP = "U+005C__DROP__"
5757
NACK = "U+005C__NACK__"
58+
FAIL = "U+005C__FAIL__"
5859

5960
_PROCESS_COUNT = os.cpu_count()
6061
# Cap max value to 16

packages/pynumaflow/pynumaflow/batchmapper/__init__.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
from pynumaflow._constants import DROP
1+
from pynumaflow._constants import DROP, FAIL
22

33
from pynumaflow.batchmapper._dtypes import (
44
Message,
@@ -14,6 +14,7 @@
1414
"Message",
1515
"Datum",
1616
"DROP",
17+
"FAIL",
1718
"BatchMapAsyncServer",
1819
"BatchMapper",
1920
"BatchResponses",

packages/pynumaflow/pynumaflow/batchmapper/_dtypes.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@
55
from typing import TypeAlias, TypeVar
66
from collections.abc import AsyncIterable, Callable
77

8-
from pynumaflow._constants import DROP, NACK
8+
from pynumaflow._constants import DROP, FAIL, NACK
99
from pynumaflow._nack import NackOptions
1010
from pynumaflow._validate import _validate_message_fields
1111

@@ -51,6 +51,10 @@ def to_nack(cls: type[M], opts: NackOptions | None = None) -> M:
5151
m._nack_options = opts
5252
return m
5353

54+
@classmethod
55+
def to_fail(cls: type[M]) -> M:
56+
return cls(b"", None, [FAIL])
57+
5458
@property
5559
def nack_options(self) -> NackOptions | None:
5660
return self._nack_options

packages/pynumaflow/pynumaflow/mapper/__init__.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,12 +5,14 @@
55
from pynumaflow.mapper._dtypes import Message, Messages, Datum, DROP, Mapper
66
from pynumaflow._metadata import UserMetadata, SystemMetadata
77
from pynumaflow._nack import NackOptions
8+
from pynumaflow._constants import FAIL
89

910
__all__ = [
1011
"Message",
1112
"Messages",
1213
"Datum",
1314
"DROP",
15+
"FAIL",
1416
"Mapper",
1517
"MapServer",
1618
"MapAsyncServer",

packages/pynumaflow/pynumaflow/mapper/_dtypes.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
from collections.abc import Callable
77
from warnings import warn
88

9-
from pynumaflow._constants import DROP, NACK
9+
from pynumaflow._constants import DROP, FAIL, NACK
1010
from pynumaflow._nack import NackOptions
1111
from pynumaflow._metadata import UserMetadata, SystemMetadata
1212
from pynumaflow._validate import _validate_message_fields
@@ -61,6 +61,10 @@ def to_nack(cls: type[M], opts: NackOptions | None = None) -> M:
6161
m._nack_options = opts
6262
return m
6363

64+
@classmethod
65+
def to_fail(cls: type[M]) -> M:
66+
return cls(b"", None, [FAIL])
67+
6468
@property
6569
def nack_options(self) -> NackOptions | None:
6670
return self._nack_options

packages/pynumaflow/pynumaflow/mapstreamer/__init__.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
from pynumaflow._constants import DROP
1+
from pynumaflow._constants import DROP, FAIL
22

33
from pynumaflow.mapstreamer._dtypes import Message, Messages, Datum, MapStreamer
44
from pynumaflow.mapstreamer.async_server import MapStreamAsyncServer
@@ -9,6 +9,7 @@
99
"Messages",
1010
"Datum",
1111
"DROP",
12+
"FAIL",
1213
"MapStreamAsyncServer",
1314
"MapStreamer",
1415
"NackOptions",

packages/pynumaflow/pynumaflow/mapstreamer/_dtypes.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
from collections.abc import AsyncIterable, Callable
77
from warnings import warn
88

9-
from pynumaflow._constants import DROP, NACK
9+
from pynumaflow._constants import DROP, FAIL, NACK
1010
from pynumaflow._nack import NackOptions
1111
from pynumaflow._validate import _validate_message_fields
1212

@@ -51,6 +51,10 @@ def to_nack(cls: type[M], opts: NackOptions | None = None) -> M:
5151
m._nack_options = opts
5252
return m
5353

54+
@classmethod
55+
def to_fail(cls: type[M]) -> M:
56+
return cls(b"", None, [FAIL])
57+
5458
@property
5559
def nack_options(self) -> NackOptions | None:
5660
return self._nack_options

packages/pynumaflow/pynumaflow/sourcetransformer/__init__.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,12 +10,14 @@
1010
from pynumaflow.sourcetransformer.async_server import SourceTransformAsyncServer
1111
from pynumaflow._metadata import UserMetadata, SystemMetadata
1212
from pynumaflow._nack import NackOptions
13+
from pynumaflow._constants import FAIL
1314

1415
__all__ = [
1516
"Message",
1617
"Messages",
1718
"Datum",
1819
"DROP",
20+
"FAIL",
1921
"SourceTransformServer",
2022
"SourceTransformer",
2123
"SourceTransformMultiProcServer",

packages/pynumaflow/pynumaflow/sourcetransformer/_dtypes.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
from collections.abc import Awaitable, Callable
77
from warnings import warn
88

9-
from pynumaflow._constants import DROP, NACK
9+
from pynumaflow._constants import DROP, FAIL, NACK
1010
from pynumaflow._nack import NackOptions
1111
from pynumaflow._metadata import UserMetadata, SystemMetadata
1212
from pynumaflow._validate import _validate_message_fields
@@ -70,6 +70,10 @@ def to_nack(
7070
m._nack_options = opts
7171
return m
7272

73+
@classmethod
74+
def to_fail(cls: type[M], event_time: datetime) -> M:
75+
return cls(b"", event_time, None, [FAIL])
76+
7377
@property
7478
def nack_options(self) -> NackOptions | None:
7579
return self._nack_options

packages/pynumaflow/tests/batchmap/test_messages.py

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import pytest
22

33
from pynumaflow.batchmapper import Message, DROP, BatchResponse, BatchResponses, NackOptions
4-
from pynumaflow._constants import NACK
4+
from pynumaflow._constants import FAIL, NACK
55
from tests.batchmap.test_datatypes import TEST_ID
66
from tests.testing_utils import mock_message
77

@@ -32,6 +32,14 @@ def test_message_default_nack_options():
3232
assert msg.nack_options is None
3333

3434

35+
def test_message_to_fail():
36+
msg = Message.to_fail()
37+
assert type(msg) is Message
38+
assert msg.keys == []
39+
assert msg.value == b""
40+
assert msg.tags == [FAIL]
41+
42+
3543
def test_batch_responses_init():
3644
batch_responses = BatchResponses()
3745
batch_response1 = BatchResponse.from_id(TEST_ID)

0 commit comments

Comments
 (0)