diff --git a/tests/kafka_test.py b/tests/kafka_test.py index cfd3fde..50f8faf 100644 --- a/tests/kafka_test.py +++ b/tests/kafka_test.py @@ -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()]) diff --git a/utils_v2/queue/kafka/models/message.py b/utils_v2/queue/kafka/models/message.py index 80e7029..a706d6d 100644 --- a/utils_v2/queue/kafka/models/message.py +++ b/utils_v2/queue/kafka/models/message.py @@ -39,9 +39,9 @@ sys.path.append("..") from pydantic import BaseModel, Field, field_validator, model_validator, AwareDatetime from typing import Optional, Literal, Union, Dict, List, Any -# Related to Google: -from google.auth.transport.requests import Request -from google.oauth2.credentials import Credentials +# # Related to Google: +# from google.auth.transport.requests import Request +# from google.oauth2.credentials import Credentials # My utils: from utils_v2.string import json