数据预处理-数据推送-代码实现|学习笔记

简介: 快速学习数据预处理-数据推送-代码实现

开发者学堂课程【大数据实战项目:反爬虫系统(Lua+Spark+Redis+Hadoop框架搭建)第四阶段数据预处理-数据推送-代码实现】学习笔记,与课程紧密联系,让用户快速学习知识。

课程地址:https://developer.aliyun.com/learning/course/672/detail/11671


数据预处理-数据推送-代码实现


数据推送-代码实现

已经过滤出纯查询的数据,将查询的数据推送到查询的 topic 中,数据已具备,要将数据推送到 Kafka 中,先读取配置文件中的配置,查询的 topic 已经配置好,在提供的 Kafka 文件中,打开文件

#消费者

#来自采集服务的原始数据

source.nginx.topic = B2CDATA_COLLECTION3

#处理后的查询数据

source.query .topic = processedouery

#处理后的预订数据

source.book .topic = processedBook

#生产者

#推送查询数据

target.query.topic = processedouery

#推送预订数据

target.book.topic = processedBook

采集完数据推送到 Kafka,即流程中的第二步

image.png

要完成第四步,即拿到 topic

生产者、推送查询数据中的 topic,将查询的数据推送到查询的 topic,即target.query.topic = processedouery,将预定的数据推送到预定的 topic,即 target.book.topic = processedBook

读取数据,使用 PropertiesUtil 调用,key有两个值,一个是推送查询数据的 topic,第二个参数是配置文件名称,即kafkaConfig.properties,推送查询数据的 topic 拿到,定义查询变量 queryTopic

创建Kafka生产者,首先拿到数据,遍历数据分区,需要遍历多个 partition,foreachPartition 效率比 foreach 快,先遍历分区,在一个分区创建生产者,一个分区创建一个生产者,多个分区有多个生产者,多个生产者同时写出速度更快,效率更高,创建一个 KafkaProducer 变量等于新的 KafkaProducer,范型是 string,需要一个Kafka参数,实例Kafka参数,val props=new,用map类型进行封装,定义java类型的 util.HashMap,map 有k和v,k是 string类型,v是 object 类型,实现参数往里面添加数据put,k指定 Kafka 集群

default.brokers = 192.168.100.100:9092,192.168.100.110:9092,192.168.100.120:9092

作为v添加

k 调用 producerConfig,引用 org.apache,

image.png

BOOTSTRAP_SERVERS 属性,将集群配置文件的值加入,k 是 default.brokers,v 是 kafkaConfig.properties,Kafka 配置文件名称

Key 的序列化、value 的序列化以及一个批次提交数据大小或间隔的时间都要进行配置

依次将配置文件名称修改,key发生变化,v不需要改变,就是 Kafka 配置文件名称,更改 ProducerConfig 的配置,使得前后一致

引用 org.apache,BOOTSTRAP_SERVERS 属性是因为直接设置好,可以直接使用

配置文件引用完成后,直接上传到生产者的参数中,流程与写Kafka的流程是一样的,数据生产者引用完成,下一步数据的载体

Partition 是分区,载体要拿到一条条数据,遍历分区数据,Partition 直接调用 foreach 或 map 数据就能拿到一个个的结果,需要返回值用 map,不需要返回值用 foreach,这里不需要返回值直接使用 foreach,foreach 拿到每一个数据 message,遍历出某一条数据,一个数据一个载体,进行下一步数据的载体,定义变量 record 等于新的ProducerRecord,需要一个 string 类型的范型,传入 queryTopic 中,将数据 message 写入到 Topic 中,就是数据的载体,数据载体具备后,发送数据,用生产者 KafkaProducer,send(record),数据发送完关闭生产者,

//将数据推送到 kafka

// 1在配置文件中读取查询类的Topic到程序中

val queryTopic= propertiesutil.getstringByKey( key = "target. query.topic" , propName ="kafkaConfig.properties")

//实例 kafka 参数

val props=new util.HashMap[string,object]()

props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFI6,Propertiesutil.getstringByKey(key = "default.brokers",propName = "kafkaconfig.properties"))props.put(ProducerConfig.KEY_SERIALIZER_CLASs_CONFIG,Propertiesutil.getstringByKey(key = "default.key_serializer_class_config",propName = “kafkaconfig.properties")) props.put(Producercojfig.vALUE_SERIALIZER_CLASs_CONFTG,PropertiesUtil.getstringByKey(key = "default.value_serializer_class_config", propName = "kafkaConfig.properties"))props.put(Producerconfig.BATCH_SIZE_CONFTG ,PropertiesUtil.getstringByKey(key = "default.batch_size_config",propName = "kafkaConfig.properties")) props.put(ProducerConfig.LINGER_MS_CONFIG,Propertiesutil.getstringByKey( key = "default.linger_ms_config",propName = "kafkaconfig.properties"))

//遍历数据的分区

queryDatas.foreachPartition(partition=>{

// 2创建kafka生产者

val kafkaProducer=new KafkaProducer[string,string](props)

//遍历partition 内的数据

partition.foreach(message=>{

//3数据的载体

val record=new ProducerRecord[string,string](queryTopic,message)

//4数据的发送

kafkaProducer.send(record)

})

//5关闭生成者

kafkaProducer.close()

推送数据到 Kafka 的过程已写完,回到程序调用的方法中,做一个接收,

//9数据推送

//9-1查询类数据的推送

Val Datasend.sendQueryDataToKafka(DataProcess)

相关文章
|
Web App开发 索引
流媒体服务器SRS部署
github地址:https://github.com/ossrs/srs 1,srs下载 http://ossrs.net/srs.release/releases/index.html 选择正式发形版 2,安装 # unzip SRS-CentOS6-x86_64-1.
5039 0
|
8月前
|
数据采集 人工智能 搜索推荐
AI 问答占 52%!长沙别墅装修 GEO 突围:30 天引用率暴涨 40%
周有贵,巴黎学院人工智能博士,GGI商学院GEO首席技术专家,专注AI时代数字营销革新。2025年12月1日,长沙著名别墅设计师张主华专程拜访交流,共探GEO技术在装修设计行业中的AI引流逻辑与实操应用。面对生成式AI问答入口占比突破52%的新趋势,传统SEO正被GEO取代——从链接点击到答案呈现,企业需通过构建灯塔内容、E-E-A-T信任链与结构化数据,让品牌信息被AI优先引用。本次对话揭示:未来流量之争,本质是“被AI推荐”的能力之争。
|
9月前
|
人工智能
1688新灯塔体系全面升级!深度解析三大巨变与商家应对策略
11月7日,1688平台正式发布了“新灯塔”服务体系的全新升级公告。此次升级并非微调,而是从指标模型、权重分配到考核精度的一次系统性重构,旨在更精准地衡量与驱动商家服务能力的提升。本文将为您深度剖析此次升级的核心变化,并为企业提供切实可行的运营建议。
|
存储 安全 API
权限申请被拒?详解京东/淘宝API审核标准与申诉技巧
在对接电商 API 时,权限申请常因资质或材料问题被拒。本文详解京东、淘宝的审核标准与申诉策略,结合实战案例,教你如何提升通过率,规避风险,实现高效对接。
|
存储 安全 物联网
什么是安全密钥,它是如何工作的
安全密钥是一种物理设备,常用于双因素或多因素身份验证(2FA/MFA),以提升在线账户安全性。它通过公钥加密协议(如FIDO U2F/FIDO2)实现强大的防网络钓鱼和凭证盗窃功能。常见的类型包括USB-A、USB-C、NFC和蓝牙密钥,支持一键登录且兼容多种服务。即使凭据泄露,安全密钥也能有效保护账户。若丢失密钥,可通过备用验证码或替代验证方法恢复访问,并重新注册新密钥。工具如ADSelfService Plus可与安全密钥无缝集成,提供自适应MFA及密码管理功能,增强整体安全性。
1545 0
什么是安全密钥,它是如何工作的
|
机器学习/深度学习 计算机视觉
《深度剖析:残差连接如何攻克深度卷积神经网络的梯度与退化难题》
残差连接通过引入“短路”连接,解决了深度卷积神经网络(CNN)中随层数增加而出现的梯度消失和退化问题。它使网络学习输入与输出之间的残差,而非直接映射,从而加速训练、提高性能,并允许网络学习更复杂的特征。这一设计显著提升了深度学习在图像识别等领域的应用效果。
918 13
|
数据采集 缓存 前端开发
服务器端渲染(SSR)
服务器端渲染(SSR)
|
缓存 安全 网络协议
HTTP中如何正确使用Via
【10月更文挑战第20天】Via`首部字段记录报文途中每个代理或网关信息,助于诊断问题和避免循环。
|
SQL 数据库 数据库管理
逆天了!IDEA执行大文件SQL,效率甩 Navicat 几条街?
【10月更文挑战第1天】在数据库管理和开发领域,SQL文件的执行效率是衡量数据库管理工具性能的重要指标之一。近期,IDEA(IntelliJ IDEA)在执行大文件SQL方面的表现引起了广泛关注,其效率远超传统的数据库管理工具Navicat。本文将深入探讨这一现象背后的原因,并结合工作学习中的技术干货,为大家带来一些实用的建议和技巧。
688 1
|
Java jenkins 测试技术
云效Flow:打造高效、稳定的CI/CD流程实战指南
云效流水线Flow评测展示新建流水线步骤,包括选择模板、添加源、Java构建、主机部署及自定义任务。通过图形界面逐项配置,如代码扫描,保存后运行流水线。虽然Flow易于上手,功能丰富,支持多环境部署,但复杂项目管理稍显繁琐,社区支持需加强。对比其他CI/CD工具,Flow在成本、功能和性能上有竞争力,适合作为团队选择。
云效Flow:打造高效、稳定的CI/CD流程实战指南