教学文库网 - 权威文档分享云平台
您的当前位置:首页 > 精品文档 > 高等教育 >

Spark整合kafka0.10.0新特性(一)(5)

来源:网络收集 时间:2026-09-11
导读: def Subscribe[K, V]( topics: ju.Collection[jl.String], kafkaParams: ju.Map[String, Object], offsets: ju.Map[TopicPartition, jl.Long]): ConsumerStrategy[K, V] = { new Subscribe[K, V](topics, kafkaPara

def Subscribe[K, V](

topics: ju.Collection[jl.String],

kafkaParams: ju.Map[String, Object],

offsets: ju.Map[TopicPartition, jl.Long]): ConsumerStrategy[K, V] = {

new Subscribe[K, V](topics, kafkaParams, offsets) } /**

* :: Experimental ::

* Subscribe to a collection of topics.

* @param topics collection of topics to subscribe * @param kafkaParams Kafka

* configuration parameters to be used on driver. The same params will be used on executors,

* with minor automatic modifications applied. * Requires \

* with Kafka broker(s) specified in host1:port1,host2:port2 form. */

@Experimental

def Subscribe[K, V](

topics: ju.Collection[jl.String],

kafkaParams: ju.Map[String, Object]): ConsumerStrategy[K, V] = { new Subscribe[K, V](topics, kafkaParams, ju.Collections.emptyMap[TopicPartition, jl.Long]())

}

/** :: Experimental ::

* Subscribe to all topics matching specified pattern to get dynamically assigned partitions. * The pattern matching will be done periodically against topics existing at the time of check. * @param pattern pattern to subscribe to * @param kafkaParams Kafka

*

* configuration parameters to be used on driver. The same params will be used on executors,

* with minor automatic modifications applied. * Requires \

* with Kafka broker(s) specified in host1:port1,host2:port2 form.

* @param offsets: offsets to begin at on initial startup. If no offset is given for a * TopicPartition, the committed offset (if applicable) or kafka param * auto.offset.reset will be used. */

@Experimental

def SubscribePattern[K, V](

pattern: ju.regex.Pattern,

kafkaParams: collection.Map[String, Object], offsets: collection.Map[TopicPartition, Long]): ConsumerStrategy[K, V] = {

new SubscribePattern[K, V]( pattern,

new ju.HashMap[String, Object](kafkaParams.asJava),

new ju.HashMap[TopicPartition, jl.Long](offsets.mapValues(l => new jl.Long(l)).asJava)) }

/** :: Experimental ::

* Subscribe to all topics matching specified pattern to get dynamically assigned partitions. * The pattern matching will be done periodically against topics existing at the time of check. * @param pattern pattern to subscribe to * @param kafkaParams Kafka

configuration parameters to be used on driver. The same params will be used on executors, * with minor automatic modifications applied. * Requires \

* with Kafka broker(s) specified in host1:port1,host2:port2 form. */

@Experimental

def SubscribePattern[K, V](

pattern: ju.regex.Pattern, kafkaParams: collection.Map[String, Object]): ConsumerStrategy[K, V] = {

new SubscribePattern[K, V]( pattern,

new ju.HashMap[String, Object](kafkaParams.asJava), ju.Collections.emptyMap[TopicPartition, jl.Long]()) }

/** :: Experimental ::

* Subscribe to all topics matching specified pattern to get dynamically assigned partitions. * The pattern matching will be done periodically against topics existing at the time of check. * @param pattern pattern to subscribe to * @param kafkaParams Kafka

* configuration parameters to be used on driver. The same params will be used on executors,

* with minor automatic modifications applied. * Requires \

* with Kafka broker(s) specified in host1:port1,host2:port2 form.

* @param offsets: offsets to begin at on initial startup. If no offset is given for a * TopicPartition, the committed offset (if applicable) or kafka param * auto.offset.reset will be used. */

@Experimental

def SubscribePattern[K, V](

pattern: ju.regex.Pattern,

kafkaParams: ju.Map[String, Object], offsets: ju.Map[TopicPartition, jl.Long]): ConsumerStrategy[K, V] = {

new SubscribePattern[K, V](pattern, kafkaParams, offsets) }

/** :: Experimental ::

* Subscribe to all topics matching specified pattern to get dynamically assigned partitions. * The pattern matching will be done periodically against topics existing at the time of check. * @param pattern pattern to subscribe to * @param kafkaParams Kafka

* configuration parameters to be used on driver. The same params will be used on executors,

* with minor automatic modifications applied. * Requires \

* with Kafka broker(s) specified in host1:port1,host2:port2 form. */

@Experimental

def SubscribePattern[K, V](

pattern: ju.regex.Pattern,

kafkaParams: ju.Map[String, Object]): ConsumerStrategy[K,

V] = {

new SubscribePattern[K, V]( pattern,

kafkaParams,

ju.Collections.emptyMap[TopicPartition, jl.Long]()) } /**

* :: Experimental ::

* Assign a fixed collection of TopicPartitions

* @param topicPartitions collection of TopicPartitions to assign * @param kafkaParams Kafka

*

* configuration parameters to be used on driver. The same params will be used on executors,

* with min …… 此处隐藏:3621字,全部文档内容请下载后查看。喜欢就下载吧 ……

Spark整合kafka0.10.0新特性(一)(5).doc 将本文的Word文档下载到电脑,方便复制、编辑、收藏和打印
本文链接:https://www.jiaowen.net/wendang/613066.html(转载请注明文章来源)
Copyright © 2020-2025 教文网 版权所有
声明 :本网站尊重并保护知识产权,根据《信息网络传播权保护条例》,如果我们转载的作品侵犯了您的权利,请在一个月内通知我们,我们会及时删除。
客服QQ:78024566 邮箱:78024566@qq.com
苏ICP备19068818号-2
Top
× 游客快捷下载通道(下载后可以自由复制和排版)
VIP包月下载
特价:29 元/月 原价:99元
低至 0.3 元/份 每月下载150
全站内容免费自由复制
VIP包月下载
特价:29 元/月 原价:99元
低至 0.3 元/份 每月下载150
全站内容免费自由复制
注:下载文档有可能出现无法下载或内容有问题,请联系客服协助您处理。
× 常见问题(客服时间:周一到周五 9:30-18:00)