【问题标题】:test on a value of a Kafka stream position测试 Kafka 流位置的值
【发布时间】: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


    【解决方案1】:

    flag == "2" 永远不会为真,因为那是 Spark RDD 对象,而不是单数字符串。

    另外,Kafka 可能有连续的记录流,因此仅检查第一条记录的第二列(假设您调用了 collect() 函数)也不起作用。

    如果要检查任何行的 2,则必须对其进行过滤

    lines = kafkaStream.map(lambda x: x[1])
    flag = lines.map(lambda line: line.split(",")).filter(lambda columns: columns[1] == "2")
    flag.pprint()
    

    如果您只想使用 Python 使用 Kafka 并检查记录值,则不需要 Spark

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-09-25
      • 1970-01-01
      • 1970-01-01
      • 2020-04-02
      • 1970-01-01
      • 2014-11-14
      • 2019-09-23
      • 2021-09-28
      相关资源
      最近更新 更多