Uh oh!
There was an error while loading. Please reload this page.
- Notifications
You must be signed in to change notification settings - Fork 52
feat(asyncio): add Reader API#309
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Uh oh!
There was an error while loading. Please reload this page.
Changes from all commits
File filter
Filter by extension
Conversations
Uh oh!
There was an error while loading. Please reload this page.
Jump to
Uh oh!
There was an error while loading. Please reload this page.
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -17,6 +17,7 @@ | ||
| * under the License. | ||
| */ | ||
| #include "utils.h" | ||
| #include <pybind11/functional.h> | ||
| #include <pybind11/pybind11.h> | ||
| namespace py = pybind11; | ||
| @@ -54,16 +55,46 @@ void Reader_seek_timestamp(Reader& reader, uint64_t timestamp) { | ||
| bool Reader_is_connected(Reader& reader) { return reader.isConnected(); } | ||
| void Reader_readNextAsync(Reader& reader, ReadNextCallback callback) { | ||
| py::gil_scoped_release release; | ||
| reader.readNextAsync(callback); | ||
| } | ||
BewareMyPower marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| void Reader_closeAsync(Reader& reader, ResultCallback callback) { | ||
| py::gil_scoped_release release; | ||
| reader.closeAsync(callback); | ||
| } | ||
| void Reader_seekAsync(Reader& reader, const MessageId& msgId, ResultCallback callback) { | ||
| py::gil_scoped_release release; | ||
| reader.seekAsync(msgId, callback); | ||
| } | ||
| void Reader_seekAsync_timestamp(Reader& reader, uint64_t timestamp, ResultCallback callback) { | ||
| py::gil_scoped_release release; | ||
| reader.seekAsync(timestamp, callback); | ||
| } | ||
| void Reader_hasMessageAvailableAsync(Reader& reader, HasMessageAvailableCallback callback) { | ||
| py::gil_scoped_release release; | ||
| reader.hasMessageAvailableAsync(callback); | ||
| } | ||
| void export_reader(py::module_& m) { | ||
| using namespace py; | ||
| class_<Reader>(m, "Reader") | ||
| .def("topic", &Reader::getTopic, return_value_policy::copy) | ||
| .def("read_next", &Reader_readNext) | ||
| .def("read_next", &Reader_readNextTimeout) | ||
| .def("read_next_async", &Reader_readNextAsync) | ||
BewareMyPower marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| .def("has_message_available", &Reader_hasMessageAvailable) | ||
| .def("has_message_available_async", &Reader_hasMessageAvailableAsync) | ||
| .def("close", &Reader_close) | ||
| .def("close_async", &Reader_closeAsync) | ||
| .def("seek", &Reader_seek) | ||
| .def("seek", &Reader_seek_timestamp) | ||
| .def("seek_async", &Reader_seekAsync) | ||
| .def("seek_async", &Reader_seekAsync_timestamp) | ||
| .def("is_connected", &Reader_is_connected); | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -39,6 +39,7 @@ | ||
| Consumer, | ||
| Producer, | ||
| PulsarException, | ||
| Reader, | ||
| _set_future, | ||
| ) | ||
| from pulsar.schema import ( # pylint: disable=import-error | ||
| @@ -465,6 +466,86 @@ async def test_seek_timestamp(self): | ||
| msg = await consumer.receive() | ||
| self.assertEqual(msg.data(), b'msg-3') | ||
| async def test_reader_simple(self): | ||
BewareMyPower marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| topic = f'asyncio-test-reader-simple-{time.time()}' | ||
| reader = await self._client.create_reader(topic, pulsar.MessageId.earliest) | ||
| self.assertTrue(reader.is_connected()) | ||
| self.assertEqual(reader.topic(), f'persistent://public/default/{topic}') | ||
| producer = await self._client.create_producer(topic) | ||
| await producer.send(b'hello') | ||
| msg = await reader.read_next() | ||
| self.assertEqual(msg.data(), b'hello') | ||
| with self.assertRaises(asyncio.TimeoutError): | ||
| await asyncio.wait_for(reader.read_next(), 1) | ||
BewareMyPower marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| await reader.close() | ||
| self.assertFalse(reader.is_connected()) | ||
| async def test_reader_on_last_message(self): | ||
| topic = f'asyncio-test-reader-on-last-message-{time.time()}' | ||
| producer = await self._client.create_producer(topic) | ||
| for i in range(10): | ||
| await producer.send(f'hello-{i}'.encode()) | ||
| reader = await self._client.create_reader(topic, pulsar.MessageId.latest) | ||
| for i in range(10, 20): | ||
| await producer.send(f'hello-{i}'.encode()) | ||
| for i in range(10, 20): | ||
| msg = await reader.read_next() | ||
| self.assertEqual(msg.data(), f'hello-{i}'.encode()) | ||
| await reader.close() | ||
| async def test_reader_on_specific_message(self): | ||
| topic = f'asyncio-test-reader-on-specific-msg-{time.time()}' | ||
| producer = await self._client.create_producer(topic) | ||
| msg_ids = [] | ||
| for i in range(10): | ||
| msg_id = await producer.send(f'hello-{i}'.encode()) | ||
| msg_ids.append(msg_id) | ||
| reader1 = await self._client.create_reader(topic, pulsar.MessageId.earliest) | ||
| for i in range(5): | ||
| msg = await reader1.read_next() | ||
| self.assertEqual(msg.data(), f'hello-{i}'.encode()) | ||
| last_msg_id = msg_ids[4] | ||
| reader2 = await self._client.create_reader(topic, last_msg_id) | ||
| for i in range(5, 10): | ||
| msg = await reader2.read_next() | ||
| self.assertEqual(msg.data(), f'hello-{i}'.encode()) | ||
| await reader1.close() | ||
| await reader2.close() | ||
| async def test_reader_has_message_available(self): | ||
| topic = f'asyncio-test-reader-has-message-available-{time.time()}' | ||
| producer = await self._client.create_producer(topic) | ||
| reader = await self._client.create_reader(topic, pulsar.MessageId.latest) | ||
| self.assertFalse(await reader.has_message_available()) | ||
| for i in range(10): | ||
| await producer.send(f'hello-{i}'.encode()) | ||
| for _ in range(10): | ||
| self.assertTrue(await reader.has_message_available()) | ||
| await reader.read_next() | ||
| self.assertFalse(await reader.has_message_available()) | ||
| await reader.close() | ||
| async def test_reader_seek(self): | ||
| topic = f'asyncio-test-reader-seek-{time.time()}' | ||
| producer = await self._client.create_producer(topic) | ||
| msg_ids = [] | ||
| for i in range(10): | ||
| msg_id = await producer.send(f'msg-{i}'.encode()) | ||
| msg_ids.append(msg_id) | ||
| reader = await self._client.create_reader(topic, pulsar.MessageId.latest, | ||
| start_message_id_inclusive=False) | ||
| await reader.seek(msg_ids[2]) | ||
| msg = await reader.read_next() | ||
| self.assertEqual(msg.data(), b'msg-3') | ||
| await reader.close() | ||
| reader_inclusive = await self._client.create_reader(topic, pulsar.MessageId.latest, | ||
| start_message_id_inclusive=True) | ||
| await reader_inclusive.seek(msg_ids[2]) | ||
| msg = await reader_inclusive.read_next() | ||
| self.assertEqual(msg.data(), b'msg-2') | ||
| await reader_inclusive.close() | ||
| async def test_schema(self): | ||
| class ExampleRecord(Record): # pylint: disable=too-few-public-methods | ||
| """Example record schema for testing.""" | ||
| @@ -507,6 +588,12 @@ def raise_exception(): | ||
| self.assertEqual(e.exception.error(), pulsar.Result.AuthenticationError) | ||
| # TODO: we should fix the error message not included in pattern subscription case | ||
| with self.assertRaises(PulsarException) as e: | ||
| await client.create_reader("private/auth/asyncio-test-token-auth-reader", | ||
| pulsar.MessageId.earliest) | ||
| self.assertEqual(e.exception.error(), pulsar.Result.AuthenticationError) | ||
| self.assertIn("token supplier failed", str(e.exception)) | ||
| await client.close() | ||
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.