【发布时间】: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