首页 文章 精选 留言 我的

精选列表

搜索[网站开发],共10000篇文章
优秀的个人博客,低调大师

如何解决spark开发中遇到需要去掉文件前几行数据的问题

版权声明:本文由董可伦首发于https://dongkelun.com,非商业转载请注明作者及原创出处。商业转载请联系作者本人。 https://blog.csdn.net/dkl12/article/details/80550095 我的原创地址:https://dongkelun.com/2018/05/27/sparkDelFirstNLines/ 前言 我碰到的问题是这样的,我需要读取压缩文件里的数据存到hive表里,压缩文件解压之后是一个txt,这个txt里前几行的数据是垃圾数据,而这个txt文件太大,txt是直接打不开的,所以不能手动打开删除前几行数据,而这个文件是业务人员从别人那拿到的所以也不能改,本文就是讲如何解决这个问题。 1、数据 首先造几条数据,以理解我的需求 data.txt id name addr time ------------ ------------------- --------------- -------------------- 1 zhangsan shanghai 2018-05-25 2 zhangsan shanghai 2018-05-25 3 zhangsan shanghai 2018-05-25 4 zhangsan shanghai 2018-05-25 5 zhangsan shanghai 2018-05-25 其中前三行是我不想要的数据,第一行为空,第二行为字段名,第三行应该是为了美观单独加了一行。 2、尝试用代码解决 2.1 思路一 用zipWithIndex给rdd加上索引,索引从0开始依次次递增1 val path = "files/data.txt" val rdd = sc.textFile(path) println("分区数:" + rdd.getNumPartitions) val rdd1 = rdd.zipWithIndex() //过滤掉索引小于等于2的 val rdd2 = rdd1.filter(_._2 > 2) rdd1.foreach(println) println("**********分割线***********") rdd2.map(kv => kv._1).foreach(println) 分区数:1 (,0) (id name addr time ,1) (------------ ------------------- --------------- --------------------,2) (1 zhangsan shanghai 2018-05-25,3) (2 zhangsan shanghai 2018-05-25,4) (3 zhangsan shanghai 2018-05-25,5) (4 zhangsan shanghai 2018-05-25,6) (5 zhangsan shanghai 2018-05-25,7) **********分割线*********** 1 zhangsan shanghai 2018-05-25 2 zhangsan shanghai 2018-05-25 3 zhangsan shanghai 2018-05-25 4 zhangsan shanghai 2018-05-25 5 zhangsan shanghai 2018-05-25 将rdd的分区改为8(大于1即可)测试一下 将 val rdd = sc.textFile(path) 改为 val rdd = sc.textFile(path, 8) 发现结果是一样的 因为我不太熟悉读取本地数据分区和读取hdfs数据分区的是否一样,所以将数据放在分布式的hdfs上测试一下 首先将data.txt上传的hdfs上 hadoop fs -put data.txt /tmp/dkl/ 将代码中的path改为 val path = "hdfs://ambari.master.com:8020/tmp/dkl/data.txt" 发现最后的结果也是一样的(分区数可能不一样) 那么这样看来这个思路是可以解决这个问题的,但是我举得例子数据量比较少,在实际工作中数据量大的话,用zipWithIndex会有性能问题。 2.2 思路2 尝试直接获取rdd的前几行数据,然后过滤掉这几行数据,但是这个前提是前几行数据在rdd里是唯一的,否在会过滤掉其他行一样的数据,我的使用场景是删掉垃圾数据,如果其他行也有一样的数据,那么正好删掉了~ val path = "files/data.txt" val rdd = sc.textFile(path, 8) println("分区数:" + rdd.getNumPartitions) //前三条 val arr = rdd.take(3) //过滤掉arr里的数据 val rdd3 = rdd.filter(!arr.contains(_)) rdd3.foreach(println) 结果 分区数:8 1 zhangsan shanghai 2018-05-25 2 zhangsan shanghai 2018-05-25 3 zhangsan shanghai 2018-05-25 4 zhangsan shanghai 2018-05-25 5 zhangsan shanghai 2018-05-25 从结果看是可以解决我的问题,且和分区多少无关,大家可以试一下。 2.3 关于rdd前几行的定义 一开始我对于rdd前几行的数据的定义是有疑惑的,不知道改变rdd分区的数目,通过rdd.first或者rdd.take获取到的前几条数据是否是固定的,此次通过写代码测试,确定是和分区无关的,而且我担心测试数据量过小,将1.5G的txt放在hdfs测试,并且不指定分区数目(默认)进行测试,发现分区默认大小13,前几行数据依旧是固定,由此更加确认和rdd分区无关。 2.4 关于rdd重新分区 可通过repartition和coalesce对rdd进行重新分区,通过如下代码测试 rdd.repartition(3).take(3).foreach(println) rdd.coalesce(3).take(3).foreach(println) 通过结果得出coalesce之后顺序和之前的顺序是一样的,而repartition之后的顺序和之前不一样了,也就是如果rdd进行repartition之后的前几行和原来的前几行是不一样的,但是重新分区的数目固定的话,每次repartition之后的顺序是一样的。 注:repartition 内部实现调用的 coalesce 且为coalesce中 shuffle = true的实现 关于Spark中repartition和coalesce的使用场景,参考Spark中repartition和coalesce的区别与使用场景解析 因为我没有repartition的需求,所以可以通过2.2的代码解决我的问题,且如果有repartition的需求,可以在删掉前几行之后再repartition~ 3、通过Linux命令删除文件前几行 命令如下: cat data.txt id name addr time ------------ ------------------- --------------- -------------------- 1 zhangsan shanghai 2018-05-25 2 zhangsan shanghai 2018-05-25 3 zhangsan shanghai 2018-05-25 4 zhangsan shanghai 2018-05-25 5 zhangsan shanghai 2018-05-25 tail -n+4 data.txt > data_new.txt cat data_new.txt 1 zhangsan shanghai 2018-05-25 2 zhangsan shanghai 2018-05-25 3 zhangsan shanghai 2018-05-25 4 zhangsan shanghai 2018-05-25 5 zhangsan shanghai 2018-05-25 mv data_new.txt data.txt mv:是否覆盖"data.txt"? y cat data.txt 1 zhangsan shanghai 2018-05-25 2 zhangsan shanghai 2018-05-25 3 zhangsan shanghai 2018-05-25 4 zhangsan shanghai 2018-05-25 5 zhangsan shanghai 2018-05-25 注:其中的cat只是为了便于理解每个操作步骤之后的结果 这样就可以删除文件前三条数据了,之后压缩,上传到hdfs再用spark程序处理即可 参考:linux删除大文件的前n行 4、最后 不知道我对spark的分区的理解是否正确,如果不对的话,欢迎大家提出指正~ 附录 最后附上将此格式的txt文件读取为rdd并转为df的示例程序 package com.dkl.leanring.spark.sql import org.apache.spark.sql.Row import org.apache.spark.sql.SparkSession import org.apache.spark.sql.types.StringType import org.apache.spark.sql.types.StructField import org.apache.spark.sql.types.StructType object Demo { def main(args: Array[String]): Unit = { val spark = SparkSession.builder().appName("DelFirstNLines").master("local").getOrCreate() val sc = spark.sparkContext val path = "files/data.txt" val data = sc.textFile(path, 8) val arr = data.take(3) val rdd = data.filter(!arr.contains(_)) //第二行为列名 val colName = arr(1).split(" +") // +表示根据一个或多个空格进行分割 val schema = StructType(colName.map(fieldName => StructField(fieldName, StringType, true))) val rowRDD = rdd.map(_.split(" +")).map(p => Row(p: _*)) val df = spark.createDataFrame(rowRDD, schema) df.show() spark.stop() } } +---+--------+--------+----------+ | id| name| addr| time| +---+--------+--------+----------+ | 1|zhangsan|shanghai|2018-05-25| | 2|zhangsan|shanghai|2018-05-25| | 3|zhangsan|shanghai|2018-05-25| | 4|zhangsan|shanghai|2018-05-25| | 5|zhangsan|shanghai|2018-05-25| +---+--------+--------+----------+

优秀的个人博客,低调大师

研究人员开发智能衣架Take Off,自带检测功能可对衣服进行杀菌和清洁

Take Off采用木质机身,通过5个可自动翻牌卡片做为“屏幕”来显示时间、温度和穿衣信息。 近来连续的阴雨天气,衣服总是干不了。面对这样的雨天,怎么办呢?如何晒衣服才能干的更快?近日,外国研究人员研发了一款名为Take Off的智能衣架,其采用木质机身,通过5个可自动翻牌卡片做为“屏幕”来显示时间、温度和穿衣信息。 当用户把衣物挂在Take Off顶端,Take Off就会自动开启对应相关衣物的清洁模式,包括清洁时间、温度等,然后喷出蒸汽清洁衣服,同时也会释放出负离子,清除衣服上的臭味和细菌。此外,Take Off还可以自动获取天气数据,如果是阴雨天或是雨雪天,它会提醒你多加衣服。当然,她还会告诉你室外温度,让你提前做好准备。 这款智能衣架在日常使用过程中,还考虑到了噪音问题,贴心的进行了降噪设计,日常使用中不会出现杂音。此外,Take Off还配有一个手持式吸尘器和四种不同质料的刷头,可清洁衣服上的灰尘。 研发人员表示,未来他们还会给Take Off添加更多功能,比如人脸识别或是语音交互功能。其最终目标是让Take Off不仅仅成为一款智能衣架,更是人们生活出行的好伙伴。据悉,这款智能衣架已经开启了预售,大约需要99美元。 如今,智能产品已经智能产品已经与我们的生活越来越近,Take Off的出现改变我们对衣架的看法。但相信随着科技的进步,更多的智能产品将对我们生活带来重大改变。 原文发布时间: 2018-05-09 11:40 本文作者: 星星 本文来自云栖社区合作伙伴镁客网,了解相关信息可以关注镁客网。

优秀的个人博客,低调大师

Hive计算时count sum partition by等方法在数据开发时的一些用法

本篇文章用于记录平时在做hive计算写sql时的心得epr代表字段 1、常用coalesce(epr,0)方法,可以防止当前字段为空,可以在计算时给个默认值,nvl()也可以 2、常用round(epr, 2)方法,数仓有时候数据类型为float、double,计算时会有精度问题,此方法可以用来保留位数 3、(CASE WHEN epr1 in (2,2) THEN epr2 ELSE -1 END),这个可以根据一个字段的值来定义另一个字段的值 4、epr3,SUM(CASE WHEN epr1 in (2,2) THEN epr2 ELSE -epr2 END) OVER (PARTITION BY epr3) 这种用法可以解决根据epr3聚合的字段,可以根据epr1的值来决定聚合函数里的正负号,PARTITION BY是可以解决在查询的时候可以直接聚合数据,而不需要单独group by数据 5、count(DISTINCT epr1) 对该字段去重去null的计数,count(epr1) 对该字段去null的计数 6、row_number() OVER (partition BY epr1, epr2 ORDER BY epr3 DESC) as number 先对epr1、epr2两个字段聚合数据然后在按epr3排序,按自然数顺序往下排,epr3相同比较数据的顺序,递增 7、rank() OVER (partition BY epr1, epr2 ORDER BY epr3 DESC) as number 先对epr1、epr2两个字段聚合数据然后在按epr3排序,按自然数顺序往下排,epr3相同的话rank值一样,有相等 8、hive里面group by的时候查询出的字段只能是group by 后的字段,不知道是不是我司的问题

优秀的个人博客,低调大师

分布式系统开发工具包 —— 基于Hessian的HTTP RPC调用技术

Hessian官网:http://hessian.caucho.com/ hessian是二进制web service协议。 Hessian介绍 创建Hessian服务包括四个步骤: 创建Java接口,用于提供公开服务 使用HessianProxyFactory创建客户端 创建服务实现类 在servlet引擎中配置服务 HelloWorld服务 public interface BasicAPI { public String hello(); } 服务实现 public class BasicService extends HessianServlet implements BasicAPI { private String _greeting = "Hello, world"; public void setGreeting(String greeting) { _greeting = greeting; } public String hello() { return _greeting; } } 客户端实现 String url = "http://hessian.caucho.com/test/test"; HessianProxyFactory factory = new HessianProxyFactory(); BasicAPI basic = (BasicAPI) factory.create(BasicAPI.class, url); System.out.println("hello(): " + basic.hello()); 部署标准web.xml <web-app> <servlet> <servlet-name>hello</servlet-name> <servlet-class>com.caucho.hessian.server.HessianServlet</servlet-class> <init-param> <param-name>home-class</param-name> <param-value>example.BasicService</param-value> </init-param> <init-param> <param-name>home-api</param-name> <param-value>example.Basic</param-value> </init-param> </servlet> <servlet-mapping> <url-pattern>/hello</url-pattern> <servlet-name>hello</servlet-name> </servlet-mapping> </web-app> Hessian序列化 Hessian类可以用来做序列化与反序列化 序列化 Object obj = ...; OutputStream os = new FileOutputStream("test.xml"); Hessian2Output out = new Hessian2Output(os); out.writeObject(obj); os.close(); 反序列化 InputStream is = new FileInputStream("test.xml"); Hessian2Input in = new Hessian2Input(is); Object obj = in.readObject(null); is.close(); 如果要序列化比基础类型或String类型更加复杂的java对象,务必确保对象实现了java.io.Serializable接口。 Hessian处理大量数据 分布式应用需要发送大量二进制数据时,使用InputStream会更加有效率,因为它避免了分配大量byte数组。方法参数中只有最后一个参数可能是InputStream,因为数据是在调用过程中读的。 下面是一个上传文件的API的例子 package example; public interface Upload { public void upload(String filename, InputStream data); } 如果返回结果是InputStream,客户端必须在finally块中调用InputStream.close()方法,因为Hessian不会关闭底层HTTP流,直到所有数据被读取并且input stream被关闭。 文件下载API: package example; public interface Download { public InputStream download(String filename, InputStream data); } 文件下载实现: InputStream is = fileProxy.download("test.xml"); try { ... // read data here } finally { is.close(); } 原文发布于:http://www.yesdata.net/2018/03/11/hessian/

优秀的个人博客,低调大师

开发者的2018】GAN、AutoML、统一框架、语音等十大趋势

GAN与造假 虽然生成对抗网络几年前就出现了,我对它是相当怀疑的。几年过去了,即使看到GAN在生成64x64分辨率的图像方面取得了巨大的进步,我对它仍是怀疑的。在阅读了一些数学文章之后,我更加怀疑了,因为这些文章说GAN并没有真正了解分布。但在2017年,事情有所改变。首先,一些新的有趣的架构(例如CycleGAN)和数学上改进的架构(例如Wasserstein GAN)让我实践了一些GAN网络,它们的表现一般,但在完成这两个程序之后,我确信我们可以,并且应该使用GAN来生成东西。 首先,我非常喜欢NVIDIA的一篇关于生成全高清图像的研究论文,生成的图像看起来非常真实(与一年前的64x64分辨率的令人毛骨悚然的人脸相比): 还有很多GAN在游戏行业的应用,例如用GAN生成游戏场景,英雄乃至整个世界。而且我认为我们必须意识到全新的造假水

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册