标签:spark streaming flume kafka direct java
王家林老师的课程:2016年大数据Spark“蘑菇云”行动之spark streaming消费flume采集的kafka数据Directf方式作业。
一、基本背景
Spark-Streaming获取kafka数据的两种方式Receiver与Direct的方式,本文介绍Direct的方式。具体的流程是这样的:
1、Direct方式是直接连接到kafka的节点上获取数据了。
2、基于Direct的方式:周期性地查询Kafka,来获得每个topic+partition的最新的offset,从而定义每个batch的offset的范围。
3、当处理数据的job启动时,就会使用Kafka的简单consumer api来获取Kafka指定offset范围的数据。
这种方式有如下优点:
1、简化并行读取:如果要读取多个partition,不需要创建多个输入DStream然后对它们进行union操作。Spark会创建跟Kafka partition一样多的RDD partition,并且会并行从Kafka中读取数据。所以在Kafka partition和RDD partition之间,有一个一对一的映射关系。;
2、高性能:不需要开启WAL机制,只要Kafka中作了数据的复制,那么就可以通过Kafka的副本进行恢复;
3、一次且仅一次的事务机制:Spark Streaming自己就负责追踪消费的offset,并保存在checkpoint中。
Spark自己一定是同步的,因此可以保证数据是消费一次且仅消费一次。
二、配置文件及编码
flume版本:1.6.0,此版本直接支持到kafka,不用在单独安装插件。
kafka版本2.10-0.8.2.1,必须是0.8.2.1,刚开始我用的是0.10,结果出现了下
四、各类错误大全的第2个错误。
spark版本:1.6.1。
kafka配文件:producer.properties,红色文字为特别要注意的配置坑,呵呵
#agentsection
producer.sources= s
producer.channels= c
producer.sinks= r
#sourcesection
producer.sources.s.type= exec
producer.sources.s.command= tail -f -n+1 /opt/test/test.log
producer.sources.s.channels= c
# Eachsink‘s type must be defined
producer.sinks.r.type= org.apache.flume.plugins.KafkaSink
producer.sinks.r.metadata.broker.list=192.168.0.10:9092
producer.sinks.r.partition.key=0
producer.sinks.r.partitioner.class=org.apache.flume.plugins.SinglePartition
producer.sinks.r.serializer.class=kafka.serializer.StringEncoder
producer.sinks.r.request.required.acks=0
producer.sinks.r.max.message.size=1000000
producer.sinks.r.producer.type=sync
producer.sinks.r.custom.encoding=UTF-8
producer.sinks.r.custom.topic.name=flume2kafka2streaming930
#Specifythe channel the sink should use
producer.sinks.r.channel= c
# Eachchannel‘s type is defined.
producer.channels.c.type= memory
producer.channels.c.capacity= 1000
producer.channels.c.transactionCapacity= 100
核心代码如下:
SparkConf conf = SparkConf().setMaster(). setAppName() .setJars(String[] { })Map<StringString> kafkaParameters = HashMap<StringString>()kafkaParameters.put()Set<String> topics = HashSet<String>()topics.add()JavaPairInputDStream<StringString> lines = KafkaUtils.(jscString.String.StringDecoder.StringDecoder.kafkaParameterstopics)JavaDStream<String> words = lines.flatMap(FlatMapFunction<Tuple2<StringString>String>() { Iterable<String> (Tuple2<StringString> tuple) Exception { Arrays.(tuple..split())} })JavaPairDStream<StringInteger> pairs = words.mapToPair(PairFunction<StringStringInteger>() { Tuple2<StringInteger> (String word) Exception { Tuple2<StringInteger>(word)} })JavaPairDStream<StringInteger> wordsCount = pairs.reduceByKey(Function2<IntegerIntegerInteger>() { Integer (Integer v1Integer v2) Exception { v1 + v2} })wordsCount.print()jsc.start()jsc.awaitTermination()jsc.close()
三、启动脚本
启动zookeeper
bin/zookeeper-server-start.sh config/zookeeper.properties &
启动kafka broker
bin/kafka-server-start.sh config/server.properties &
创建topic
bin/kafka-topics.sh --create --zookeeper 192.168.0.10:2181 --replication-factor 1 --partitions 1 --topic flume2kafka2streaming930
启动flume
bin/flume-ng agent --conf conf/ -f conf/producer.properties -n producer -Dflume.root.logger=INFO,console
bin/spark-submit --class com.dt.spark.sparkstreaming.SparkStreamingOnKafkaDirected --jars /lib/kafka_2.10-0.8.2.1/kafka-clients-0.8.2.1.jar,/lib/kafka_2.10-
0.8.2.1/kafka_2.10-0.8.2.1.jar,/lib/kafka_2.10-0.8.2.1/metrics-core-2.2.0.jar,/lib/spark-1.6.1/spark-streaming-kafka_2.10-1.6.1.jar --master local[5] SparkApps.jar
echo "hadoop spark hive storm spark hadoop hdfs" >> /opt/test/test.log
echo "hive storm " >> /opt/test/test.log
echo "hdfs" >> /opt/test/test.log
echo "hadoop spark hive storm spark hadoop hdfs" >> /opt/test/test.log
输出结果如下:
* 结果如下:
* -------------------------------------------
* Time: 1475282360000 ms
* -------------------------------------------
*(spark,8)
*(storm,4)
*(hdfs,4)
*(hive,4)
*(hadoop,8)
四、各类错误大全
1、Exception in thread "main" java.lang.NoClassDefFoundError: org/apache/spark/streaming/kafka/KafkaUtils
at com.dt.spark.SparkApps.SparkStreaming.SparkStreamingOnKafkaDirected.main
一概是没有提交jar包,一概会报错,无法执行,一概在submit脚本里添加:
bin/spark-submit --class com.dt.spark.sparkstreaming.SparkStreamingOnKafkaDirected --jars /lib/kafka_2.10-0.8.2.1/kafka-clients-0.8.2.1.jar,/lib/kafka_2.10-
0.8.2.1/kafka_2.10-0.8.2.1.jar,/lib/kafka_2.10-0.8.2.1/metrics-core-2.2.0.jar,/lib/spark-1.6.1/spark-streaming-kafka_2.10-1.6.1.jar --master local[5] SparkApps.jar
2、Exception in thread "main" java.lang.ClassCastException: kafka.cluster.BrokerEndPoint cannot be cast to kafka.cluster.Broker。
上stackoverflow.com及spark官网查询,这个是因为版本不兼容引起。官网提供的版本:Spark Streaming 1.6.1 is compatible with Kafka 0.8.2.1
王家林_DT大数据梦工厂
简介: 王家林:DT大数据梦工厂创始人和首席专家.微信公众号DT_Spark .
联系邮箱18610086859@vip.126.com
电话:18610086859
微信号:18610086859
微博为:http://weibo.com/ilovepains
2016年大数据Spark“蘑菇云”行动之spark streaming消费flume采集的kafka数据Directf方式
标签:spark streaming flume kafka direct java
原文地址:http://36006798.blog.51cto.com/988282/1858272