Commit 4ec30913 authored by 张彦钊's avatar 张彦钊

change

parent 73c5e48b
...@@ -166,7 +166,7 @@ def Filter_Data(data): ...@@ -166,7 +166,7 @@ def Filter_Data(data):
def write_to_kafka(): def write_to_kafka():
producer = KafkaProducer(bootstrap_servers=["172.16.44.25:9092","172.16.44.31:9092","172.16.44.45:9092"], producer = KafkaProducer(bootstrap_servers=["172.16.44.25:9092","172.16.44.31:9092","172.16.44.45:9092"],
key_serializer=lambda k: pickle.dumps(k), value_serializer=lambda v: pickle.dumps(v)) key_serializer=lambda k: json.dumps(k), value_serializer=lambda v: json.dumps(v))
print("hajs") print("hajs")
future = producer.send(topic="test_topic", key="hello", value="world") future = producer.send(topic="test_topic", key="hello", value="world")
try: try:
...@@ -191,7 +191,7 @@ def m_decoder(s): ...@@ -191,7 +191,7 @@ def m_decoder(s):
data = msgpack.loads(s,encoding='utf-8') data = msgpack.loads(s,encoding='utf-8')
return data return data
except: except:
data = pickle.loads(s) data = json.loads(s)
return data return data
if __name__ == '__main__': if __name__ == '__main__':
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment