【发布时间】:2019-11-26 16:54:55
【问题描述】:
from kafka import KafkaProducer, errors, admin, KafkaConsumer
SERVERS = ['localhost:9092']
TEST_TOPIC = 'test-topic'
DATA = [{'A':'A'}, {'A':'A'}, {'A':'A'}]
class TestKafkaConsumer(unittest.TestCase):
@classmethod
def setUpClass(self):
self._producer = KafkaProducer(bootstrap_servers=SERVERS, value_serializer=lambda x:dumps(x).encode('utf-8'))
def _send_data(self):
for data in DATA:
print(self._producer.send(TEST_TOPIC, value=data))
def test_basic_processing(self):
self._send_data()
received = []
consumer = KafkaConsumer(TEST_TOPIC, bootstrap_servers=SERVERS)
for msg in consumer:
message = json.loads(msg.value.decode('utf-8'))
received.append(message)
if (len(received) >= len(DATA)):
self.assertEqual(received, DATA)
这应该很快就会成功,因为它只是以一种非常简单的方式将数据发送到 Kafka 代理。但是,它超时了;消费者从不阅读一条消息。如果我将消费者部分移动到不同的文件并在不同的终端窗口中运行它,则消息会立即“使用”。为什么这个单元测试中的单元测试不适用于消费者?
【问题讨论】:
-
大家找到解决方案了吗?
标签: python python-3.x unit-testing apache-kafka python-unittest