(20260623) Migrate to new server.
This commit is contained in:
+3
-2
@@ -41,7 +41,7 @@ import os
|
||||
# Utils:
|
||||
from utils_v2.string import json
|
||||
from utils_v2.date_time import date_time
|
||||
from utils_v2.queue.async_kafka import ProducerKafka, ConsumerKafka, get_ssl_context
|
||||
from utils_v2.queue.kafka.controllers.async_kafka import ProducerKafka, ConsumerKafka, get_ssl_context
|
||||
|
||||
# For async activities:
|
||||
import asyncio
|
||||
@@ -134,7 +134,8 @@ if __name__ == "__main__":
|
||||
while True:
|
||||
await asyncio.sleep(1.0)
|
||||
messages = await my_consumer.consume(count = 10, timeout = 1.0)
|
||||
for message in messages: print("MESSAGE:", json.to_string(message["value"]))
|
||||
# for message in messages: print("MESSAGE:", json.to_string(message["value"]))
|
||||
for message in messages: print("MESSAGE:", json.to_string(message.value))
|
||||
|
||||
async def main():
|
||||
await asyncio.gather(*[keep_producing(), keep_consuming()])
|
||||
|
||||
Reference in New Issue
Block a user