【问题标题】:how to connect spark streaming with cassandra?如何将火花流与 cassandra 连接起来?
【发布时间】:2016-05-01 19:16:27
【问题描述】:

我正在使用

Cassandra v2.1.12
Spark v1.4.1
Scala 2.10

cassandra 正在监听

rpc_address:127.0.1.1
rpc_port:9160

比如连接kafka和spark-streaming,在每4秒听一次kafka的同时,我有如下spark作业

sc = SparkContext(conf=conf)
stream=StreamingContext(sc,4)
map1={'topic_name':1}
kafkaStream = KafkaUtils.createStream(stream, 'localhost:2181', "name", map1)

spark-streaming 每隔 4 秒持续监听一次 kafka broker 并输出内容。

同样,我希望 Spark Streaming 听 cassandra 并每隔 4 秒输出指定表的内容

如何转换上述流式代码,使其与 cassandra 而非 kafka 一起使用?


非流式解决方案

我显然可以在无限循环中继续运行查询,但这不是真正的流式传输,对吗?

火花工作:

from __future__ import print_function
import time
import sys

from random import random
from operator import add
from pyspark.streaming import StreamingContext
from pyspark import SparkContext,SparkConf
from pyspark.sql import SQLContext
from pyspark.streaming import *

sc = SparkContext(appName="sparkcassandra")
while(True):
    time.sleep(5)
    sqlContext = SQLContext(sc)
    stream=StreamingContext(sc,4)
    lines = stream.socketTextStream("127.0.1.1", 9160)
    sqlContext.read.format("org.apache.spark.sql.cassandra")\
                 .options(table="users", keyspace="keyspace2")\
                 .load()\
                 .show()

这样跑

sudo ./bin/spark-submit --packages \
datastax:spark-cassandra-connector:1.4.1-s_2.10 \
examples/src/main/python/sparkstreaming-cassandra2.py

我得到 表格值,看起来很像

lastname|age|city|email|firstname

那么从 cassandra 中“流式传输”数据的正确方法是什么?


【问题讨论】:

    标签: apache-spark cassandra pyspark spark-streaming datastax


    【解决方案1】:

    目前,从 C* 流式传输数据的“正确方法”不是从 C* 流式传输数据 :) 相反,将消息队列(如 Kafka)放在 C* 前面并流式传输通常更有意义那。 C* 不容易支持增量表读取,尽管如果集群键基于插入时间,则可以做到这一点。

    如果您有兴趣使用 C* 作为流式源代码,请务必查看并评论 https://issues.apache.org/jira/browse/CASSANDRA-8844 变更数据捕获

    这很可能是您要查找的内容。

    如果您实际上只是想定期读取整个表并做一些事情,那么您最好只使用一个启动批处理操作的 cron 作业,因为您确实无法恢复状态。

    【讨论】:

      【解决方案2】:

      目前,Spark 1.6 不支持将 Cassandra 作为流媒体源,您必须为自己的案例实现自定义接收器 (listen to cassandra and output the contents of the specified table every say 4 seconds.)。

      请参考实施指南:

      Spark Streaming Custom Receivers

      【讨论】:

        猜你喜欢
        • 2019-08-15
        • 2016-12-09
        • 2016-02-07
        • 2019-06-01
        • 1970-01-01
        • 2016-09-06
        • 2016-10-02
        • 2018-08-13
        • 2013-08-11
        相关资源
        最近更新 更多