【发布时间】:2019-04-13 08:33:38
【问题描述】:
我想测试 Kafka 流位置的值 如果相等的值具有例如“2” 然后显示启动函数 A 否则启动函数 B
kafkaStream = KafkaUtils.createDirectStream(ssc, [topic], {
'bootstrap.servers': 'localhost:9092',
'group.id': 'video-group',
'fetch.message.max.bytes': '15728640',
'auto.offset.reset': 'largest'})
# Group ID is completely arbitrary
lines = kafkaStream.map(lambda x: x[1])
flag = lines.map(lambda line: line.split(",")).map(lambda v : v[0])
if flag == "2":
A = lines.map(lambda line: line.split(",")).map(lambda v: v[1])
A.pprint()
else:
lines.pprint()
【问题讨论】:
标签: python apache-spark pyspark apache-kafka