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

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

来源:网络收集 时间:2026-09-11
导读: * with Kafka broker(s) specified in host1:port1,host2:port2 form. */ @Experimental def Assign[K, V]( topicPartitions: Iterable[TopicPartition], kafkaParams: collection.Map[String, Object]): ConsumerS

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

@Experimental def Assign[K, V](

topicPartitions: Iterable[TopicPartition],

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

new Assign[K, V](

new ju.ArrayList(topicPartitions.asJavaCollection), new ju.HashMap[String, Object](kafkaParams.asJava), 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 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 Assign[K, V](

topicPartitions: ju.Collection[TopicPartition], kafkaParams: ju.Map[String, Object],

offsets: ju.Map[TopicPartition, jl.Long]): ConsumerStrategy[K, V] = { new Assign[K, V](topicPartitions, kafkaParams, offsets) } /**

* :: Experimental ::

* Assign a fixed collection of TopicPartitions

* @param topicPartitions collection of TopicPartitions to assign * @param kafkaParams Kanc630.comfka

*

* 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 Assign[K, V](

topicPartitions: ju.Collection[TopicPartition],

kafkaParams: ju.Map[String, Object]): ConsumerStrategy[K, V] = { new Assign[K, V]( topicPartitions, kafkaParams,

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

Spark整合kafka0.10.0新特性(一)(6).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)