• [问题求助] OBS文件下载可否支持积分兑换?
    在云速建站的obs桶中并没有设置下载次数,其次是无法实现文件交易,在我们看来,文件交易比知识付费更好,像百度文库就已经实现类似于积分兑换的方式下载文件,如果云速建站也支持积分兑换下载就行了。
  • [知识分享] MRS离线数据分析:通过Flink作业处理OBS数据
    【摘要】 MRS支持在大数据存储容量大、计算资源需要弹性扩展的场景下,用户将数据存储在OBS服务中,使用MRS集群仅做数据计算处理的存算分离模式。 本文将向您介绍如何在MRS集群中运行Flink作业来处理OBS中存储的数据。本文分享自华为云社区《【云小课】EI第47课 MRS离线数据分析-通过Flink作业处理OBS数据》,作者:Hello EI 。MRS支持在大数据存储容量大、计算资源需要弹性扩展的场景下,用户将数据存储在OBS服务中,使用MRS集群仅做数据计算处理的存算分离模式。Flink是一个批处理和流处理结合的统一计算框架,其核心是一个提供了数据分发以及并行化计算的流数据处理引擎。它的最大亮点是流处理,是业界最顶级的开源流处理引擎。本文将向您介绍如何在MRS集群中运行Flink作业来处理OBS中存储的数据。Flink最适合的应用场景是低时延的数据处理(Data Processing)场景:高并发pipeline处理数据,时延毫秒级,且兼具可靠性。在本示例中,我们使用MRS集群内置的Flink WordCount作业程序,来分析OBS文件系统中保存的源数据,以统计源数据中的单词出现次数。当然您也可以获取MRS服务样例代码工程,参考Flink开发指南开发其他Flink流作业程序。本案例基本操作流程如下所示:创建MRS集群创建并购买一个包含有Flink组件的MRS集群,详情请参见购买自定义集群。本文以购买MRS 3.1.0版本的集群为例,集群未开启Kerberos认证。在本示例中,由于我们要分析处理OBS文件系统中的数据,因此在集群的高级配置参数中要为MRS集群绑定IAM权限委托,使得集群内组件能够对接OBS并具有对应文件系统目录的操作权限。您可以直接选择系统默认的“MRS_ECS_DEFAULT_AGENCY”,也可以自行创建其他具有OBS文件系统操作权限的自定义委托。集群购买成功后,在MRS集群的任一节点内,使用omm用户安装集群客户端,具体操作可参考安装并使用集群客户端。例如客户端安装目录为“/opt/client”。准备测试数据在创建Flink作业进行数据分析前,我们需要在提前准备待分析的测试数据,并将该数据上传至OBS文件系统中。本地创建一个“mrs_flink_test.txt”文件,例如文件内容如下:This is a test demo for MRS Flink. Flink is a unified computing framework that supports both batch processing and stream processing. It provides a stream data processing engine that supports data distribution and parallel computing.在云服务列表中选择“存储 > 对象存储服务”,登录OBS管理控制台。单击“并行文件系统”,创建一个并行文件系统,并上传测试数据文件。例如创建的文件系统名称为“mrs-demo-data”,单击系统名称,在“文件”页面中,新建一个文件夹“flink”,上传测试数据至该目录中。则本示例的测试数据完整路径为“obs://mrs-demo-data/flink/mrs_flink_test.txt”。上传数据分析应用程序。使用管理台界面直接提交作业时,将已开发好的Flink应用程序jar文件也可以上传至OBS文件系统中,或者MRS集群内的HDFS文件系统中。本示例中我们使用MRS集群内置的Flink WordCount样例程序,可从MRS集群的客户端安装目录中获取,即“/opt/client/Flink/flink/examples/batch/WordCount.jar”。将“WordCount.jar”上传至“mrs-demo-data/program”目录下。创建并运行Flink作业方式1:在控制台界面在线提交作业。登录MRS管理控制台,单击MRS集群名称,进入集群详情页面。在集群详情页的“概览”页签,单击“IAM用户同步”右侧的“单击同步”进行IAM用户同步。单击“作业管理”,进入“作业管理”页签。单击“添加”,添加一个Flink作业。作业类型:Flink作业名称:自定义,例如flink_obs_test。执行程序路径:本示例使用Flink客户端的WordCount程序为例。运行程序参数:使用默认值。执行程序参数:设置应用程序的输入参数,“input”为待分析的测试数据,“output”为结果输出文件。例如本示例中,我们设置为“--input obs://mrs-demo-data/flink/mrs_flink_test.txt --output obs://mrs-demo-data/flink/output”。服务配置参数:使用默认值即可,如需手动配置作业相关参数,可参考运行Flink作业。确认作业配置信息后,单击“确定”,完成作业的新增,并等待运行完成。方式2:通过集群客户端提交作业。使用root用户登录集群客户端节点,进入客户端安装目录。su - omm cd /opt/client source bigdata_env执行以下命令验证集群是否可以访问OBS。hdfs dfs -ls obs://mrs-demo-data/flink提交Flink作业,指定源文件数据进行消费。flink run -m yarn-cluster /opt/client/Flink/flink/examples/batch/WordCount.jar --input obs://mrs-demo-data/flink/mrs_flink_test.txt --output obs://mrs-demo/data/flink/output2执行后结果类似如下:... Cluster started: Yarn cluster with application id application_1654672374562_0011 Job has been submitted with JobID a89b561de5d0298cb2ba01fbc30338bc Program execution finished Job with JobID a89b561de5d0298cb2ba01fbc30338bc has finished. Job Runtime: 1200 ms查看作业执行结果作业提交成功后,登录MRS集群的FusionInsight Manager界面,选择“集群 > 服务 > Yarn”。单击“ResourceManager WebUI”后的链接进入Yarn Web UI界面,在Applications页面查看当前Yarn作业的详细运行情况及运行日志。等待作业运行完成后,在OBS文件系统中指定的结果输出文件中可查看数据分析输出的结果。下载“output”文件到本地并打开,可查看输出的分析结果。a 3 and 2 batch 1 both 1 computing 2 data 2 demo 1 distribution 1 engine 1 flink 2 for 1 framework 1 is 2 it 1 mrs 1 parallel 1 processing 3 provides 1 stream 2 supports 2 test 1 that 2 this 1 unified 1使用集群客户端命令行提交作业时,若不指定输出目录,在作业运行界面也可直接查看数据分析结果。Job with JobID xxx has finished. Job Runtime: xxx ms Accumulator Results: - e6209f96ffa423974f8c7043821814e9 (java.util.ArrayList) [31 elements] (a,3) (and,2) (batch,1) (both,1) (computing,2) (data,2) (demo,1) (distribution,1) (engine,1) (flink,2) (for,1) (framework,1) (is,2) (it,1) (mrs,1) (parallel,1) (processing,3) (provides,1) (stream,2) (supports,2) (test,1) (that,2) (this,1) (unified,1)
  • [上云精品] 筑起云端 “免疫”屏障,让你的数据有备无患
    受疫情影响,各组织的云转型进程不断加速,然而业务持续上云却给数据安全带来重重挑战。为此,华为云携手爱数,共同为组织数据安全筑起 “免疫”屏障。一、华为云携手爱数,助力组织数据保护云转型过程中,不少企业开始使用混合云或多云,数据也将同时存在于本地数据中心与公有云中。由于网络带宽与公有云出口流量成本高昂等原因,部署在本地的传统数据保护方案,难以实现对公有云数据快速且经济的保护。对此,华为云携手爱数,重磅推出数据保护服务——爱数 AnyBackup 云备份(华为云云商店在售),为混合云架构下的业务与数据,提供高效且经济的保护,助力各组织建立数据安全屏障,全面护航数据安全。爱数数据保护服务基于任意云到任意云的技术架构,为组织多云、混合云环境中的工作负载提供经济、易用、安全的数据保护服务。无需硬件资源准备,无需繁琐的安装部署,无需复杂的后期服务运维,7分钟即可开通服务,开启数据保护之旅。爱数数据保护服务助力各组织即使在疫情期间,也能实现轻松、灵活的数据保护。二、爱数数据保护服务应用场景数字经济时代,数据已成为一种重要生产要素,对组织生存与发展具有重要意义,然而数据安全却始终面临诸多挑战。1、勒索病毒攻击导致数据被锁定近年来,勒索病毒肆虐,不仅威胁着各组织的数据与财产安全,还可能让各组织随时面临业务停摆的风险。回看曾被勒索病毒攻击的组织,大多遭受了巨额经济损失,甚至数据丢失等风险。2、人为操作导致数据丢失一方面,组织内部员工删库走人事件时有发生;另一方面,即使并非主观意愿,人为误操作导致的数据删除或损坏也难以避免。重要数据的丢失,不仅会影响正常业务的运行,也会影响组织的社会公信力甚至行业地位。3、硬件故障导致数据无法使用随着计算机的长期使用,硬件故障问题难以难免,由此导致的数据无法使用问题,也是一种极具威胁的潜在数据安全隐患。数据是组织最宝贵的无形资产,考虑到威胁数据安全的种种隐患,各组织有必要未雨绸缪,部署数据保护方案,对本地或云上的业务与数据进行保护。三、爱数数据保护服务助您“有备无患”爱数数据保护服务提供的混合云数据保护方案,可以把数据安全高效地备份至华为云 OBS 对象存储中:【本地数据备份上云】对于本地环境中的数据,可以把数据直接备份至公有云,方便、快捷地享受数据保护服务带来的安全性保障;【云中数据快速备份】对于公有云中的数据,可以选择同区域的对象存储作为备份介质,提升整体备份效率,实现更高效的数据保护。爱数数据保护服务,致力于为组织业务与数据建立无懈可击的安全屏障。作为华为云云商店的大数据基础设施服务商,爱数在深刻洞察客户需求的基础上推出爱数数据保护服务,助力客户的数字化转型之旅。未来,爱数将联合华为云以更加创新的产品与技术平台为客户提供整合、治理、洞察与保护的全域数据能力,与各行各业共创数据驱动型组织。文中提到的商品链接:爱数 AnyBackup 云备份(30天免费试用,点击体验!)撰文丨爱数、格子编辑丨格子
  • [互动交流] 【OBS产品】【并行文件系统功能】 obs BrowserJS 是否支持并行文件系统
    【功能模块】【操作步骤&问题现象】1、2、【截图信息】【日志信息】(可选,上传日志内容或者附件)
  • [互动交流] 【OBS产品】【授权用户获取桶列表功能】将权限桶读取授予其他用户,其他用户如何能查询到该桶,API请求和OBS客户端都获取不到
    【功能模块】【操作步骤&问题现象】1、2、【截图信息】【日志信息】(可选,上传日志内容或者附件)
  • [技术干货] 块存储、文件存储、对象存储原理及特性,相互比较
    块存储存储提供给应用的是一个LUN或者是一个卷,LUN和卷是面向磁盘空间的一种组织方式,上层应用要通过FC或者ISCSI协议访问SAN。SAN存储处理的是管理磁盘的问题,适用于实时读写场景,如高性能计算、企业核心集群应用、企业应用系统和开发测试等,容量TB级别,时间亚毫秒级。只能在ECS/BMS中挂载使用,不能被操作系统应用直接访问,需要格式化成文件系统进行访问。文件存储提供给应用的是一个文件系统或者是一个文件夹,上层应用通过NFS和CIFS协议进行访问,利用FTP+TFTP协议进行上传下载,此外,文件系统要维护一个目录树,适用于企业组织内部共享场景,提升办公效率和存储空间利用率(减少同类型数据复存),容量PB级别,时延3-10ms。在ECS/BMS中通过网络协议挂载使用,支持NFS和CIFS的网络协议。需要指定网络地址进行访问,也可以将网络地址映射为本地目录后进行访问。对象存储更加适合web类应用,基于URL访问地址提供一个海量的桶存储空间,能够存储各种类型的文件对象,对象存储是一个扁平架构,无需维护复杂的文件目录。无需考虑存储空间的限制,一个桶支持近乎无限大的存储空间。(适用于离线、冷数据、归档数据、作为后端存储为客户打造的离线存储系统,性价比),容量EB级别,时延10ms, 可以通过互联网或专线访问。需要指定桶地址进行访问,使用的是HTTP和HTTPS等传输协议。
  • [技术干货] obs设置标准转低频,低频转归档
    生命周期管理可适用于以下典型场景:(1)周期性上传的日志文件,可能只需要保留一个星期或一个月。到期后要删除它们。(2)某些文档在一段时间内经常访问,但是超过一定时间后便可能不再访问了。这些文档需要在一定时间后转化为低频访问存储,归档存储或者删除。生命周期管理可以按对象名前缀进行设置规则,也可以在整个桶上设置规则。生命周期管理功能支持数据从当前版本转换为低频访问存储、转换为归档存储,以及数据进行过期删除。您可以指定在对象最后一次更新后多少天,受规则影响的对象将转换为低频访问存储、归档存储或者过期并自动被OBS删除。转换为低频访问存储的时间最少设置为30天,若同时设置转换为低频访问存储和转换为归档存储,则转换为归档存储的时间要比转换为低频访问存储的时间至少长30天,例如转换为低频访问存储设置为33天,则转换为归档存储至少需要设置为63天。对象存储类别转换限制:仅支持将标准存储对象转换为低频访问存储对象,低频访问存储对象转换为标准存储对象需手动转换。仅支持将标准存储或低频访问存储对象转换为归档存储对象。如果要将归档存储对象转换为标准存储或低频访问存储对象,需要手动恢复对象,然后手动转换存储类别。
  • [互动交流] 【华为OBS产品】【创建桶功能】创建名为download和test的桶时创建失败,其余名称同时创建可以成功
    【功能模块】【操作步骤&问题现象】1、2、【截图信息】【日志信息】(可选,上传日志内容或者附件)
  • [存储类] 【对象存储服务obs】【obsutil报错】在docker中异常: /bin/sh: obsutil:not found
    【功能模块】对象存储服务obs, 根据文档在docker中安装obsutil工具, 使用obsutil时, 报错:  /bin/sh: obsutil:not found。请问要怎么处理才能在docker 中执行obsutil 指令。【操作步骤&问题现象】1、根据obs.Dockerfile 文件启动一个docker容器,在docker中下载安装obsutil,到目录  /usr/local/bin;该文件存放路径默认就在docker的环境变量里面。2、执行 obsutil config xxx ,指令报错,/bin/sh: obsutil:not found;  显示找不到该指令;【截图信息】dockerfile文件内容【日志信息】(可选,上传日志内容或者附件)
  • [互动交流] 【OBS】【创建文件夹】官方示例代码无法创建文件夹
    【功能模块】官方SDK创建文件夹报错【操作步骤&问题现象】1、使用官方SDK示例代码创建文件夹,报错【截图信息】【日志信息】(可选,上传日志内容或者附件)sdk版本:<dependency> <groupId>com.huaweicloud</groupId> <artifactId>esdk-obs-java-bundle</artifactId> <version>[3.21.11,)</version></dependency>
  • [互动交流] obs文件上传
    如何在jsp页面中,直接使用  var ObsClient = require('esdk-obs-nodejs');  目前jsp中只有jquery,做一个测试用。
  • [互动交流] OBS上传照片出现问题
    在安卓设备上使用的OBS存储服务,obsClient 实例以变量形式初始化一次存储在内存上,上传的文件为照片,代码如下:String endPoint = ""; String ak = "*"; String sk = ""; // 创建ObsClient实例 ObsClient obsClient = new ObsClient(ak, sk, endPoint); // localfile为待上传的本地文件路径,需要指定到具体的文件名 obsClient.putObject("bucketname", "objectname", new File("localfile")); // localfile2 为待上传的本地文件路径,需要指定到具体的文件名 PutObjectRequest request = new PutObjectRequest(); request.setBucketName("bucketname"); request.setObjectKey("objectname2"); request.setFile(new File("localfile2")); obsClient.putObject(request);基本都会上传成功,但是小几率出现下图问题,因为没做操作日志记录,想问下这种问题原因一般是哪里出问题了
  • [互动交流] 华为云的对象存储支持鸿蒙程序访问吗?
    看华为云对象存储的SDK,支持安卓和ios访问,请问支持华为鸿蒙的程序访问吗?有操作说明吗?谢谢
  • [技术交流] 让RDS(for MySQL)数据库的慢日志、审计日志跨空间转存OBS变得更加自动化
    【摘要】 本项目将展示如何使用华为云产品RDS(for MySQL)、OBS提供的SDK方法将数据库产生的慢日志、审计日志通过网络流的方式直接将日志文件转存至OBS中进行备份存储,项目过程中能够让您熟悉OBS和RDS的SDK,详细内容可阅读文章进行了解。《目录》背景开发环境云服务介绍方案设计方案简述方案架构图时序图代码参数指南代码实现结果反馈1、背景在业务系统的运行过程中,数据库会产生一些日志,其中包括慢日志、审计日志,有时我们需要对这些日志手动的进行下载存储,有些繁琐。为了解决手动存储的弊端,希望可以通过调用接口的方式来实现定时的自动转存储至OBS效果。2、开发环境![Snipaste_2022-06-22_11-44-55.png](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/22/1655869524841820242.png)3、云服务介绍云数据库 RDS for MySQL:MySQL是目前最受欢迎的开源数据库之一,其性能卓越,搭配LAMP(Linux + Apache + MySQL + Perl/PHP/Python),成为WEB开发的高效解决方案。 云数据库 RDS for MySQL拥有稳定可靠、安全运行、弹性伸缩、轻松管理、经济实用等特点。 华为云OBS:对象存储服务(Object Storage Service,OBS)是一个基于对象的海量存储服务,为客户提供海量、安全、高可靠、低成本的数据存储能力,使用时无需考虑容量限制,并且提供多种存储类型,满足客户各类业务场景诉求。4、方案设计i、方案简述通过调用SDK提供的接口,先获取慢日志或审计日志的下载路径,然后通过网络流上传的方式直接将链接日志下载存储至OBS指定的位置。ii、方案架构图两种日志均是转存至OBS,但是因为其获取途径不一样,在某些细节实现有一些区分,如慢日志转存OBS,因没有提供时间筛选,无法控制其增量行为,所以需要通过比对前后文件来判断是否需要上传;而审计日志则是可以通过时间筛选的方法来控制,具体实现如下:慢日志转存OBS架构图![慢日志转存OBS架构图(白底板).png](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/22/1655865621378972063.png) 慢日志转存OBS实现思路: 1、通过SDK提供的接口获取日志下载链接。 2、然后new URL(fileLink).openConnection().getContentLength()获取链接文件。 3、再通过OBS SDK提供的接口方法获取OBS位置已存储的文件。 4、如果两者文件一样,则拉取下来的网络流通过OBS客户端上传至OBS指定位置进行存储审计日志转存OBS架构图![](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/22/1655865473425779338.png) 审计日志转存OBS实现思路: 1、通过SDK提供的接口方法传入时间等参数获取某一个时间段的审计日志的ID。 2、然后通过日志ID找到我们每一个日志文件的下载链接。 3、最后通过网络流将日志文件拉取下来,通过OBS客户端上传至OBS指定位置进行存储。iii、方案时序图慢日志转存OBS时序图 ![慢日志转存OBS时序图(白底版).png](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/22/1655865644641527095.png) 审计日志转存OBS时序图 ![审计日志转存OBS时序图(白底版).png](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/22/1655865655731506215.png)5、代码参数指南##### 公共参数获取方式 1. endPoint,ak,sk,areaCode,instanceID,obsName。 1. endPoint获取方式参考:https://developer.huaweicloud.com/endpoint获取对应的endPoint值 。注意:endpoint值由很多,要找到OBS下面的endpoint。 2. ak与sk获取方式可参考:https://support.huaweicloud.com/iam_faq/iam_01_0618.html 。 3. areaCode区域码获取方式,可在这查询:https://developer.huaweicloud.com/endpoint 。 4. instanceID实例ID需要在华为云的RDS实例管理界面中获取实例ID如图所示: ![实例id.png](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/21/1655778411438620812.png) 5. obsName为OBS桶名。 ##### RDSSlowlogUploadOBS参数详情 1. endPoint,ak,sk,areaCode,instanceID,obsName获取方式请参考公共参数。##### RDSAuditlogUploadOBS参数详情 1. endPoint,ak,sk,areaCode,instanceID,obsName,startTime,endTime。 1. endPoint,ak,sk,areaCode,instanceID,obsName获取方式请参考公共参数。 2. startTime:起始时间(格式:2018-08-06T10:41:14+0800) 3. endTime:终点时间(格式:2018-08-06T10:41:14+0800) 4. offset:起始查询位置(0表示从第一个开始查询) 5. limit:查询数量(一次最多能查50个)6、代码实现实现代码可参考慢日志、审计日志转存OBS7、结果反馈慢日志转存OBS ![慢日志.png](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/21/1655778861996113650.png) 审计日志转存OBS ![审计日志.png](https://bbs-img.huaweicloud.com/data/forums/attachment/forum/20226/21/1655778872989527175.png)
  • [最佳实践] 使用Spark Jar作业读取和查询OBS数据
    操作场景DLI完全兼容开源的Apache Spark,支持用户开发应用程序代码来进行作业数据的导入、查询以及分析处理。本示例从编写Spark程序代码读取和查询OBS数据、编译打包到提交Spark Jar作业等完整的操作步骤说明来帮助您在DLI上进行作业开发。环境准备在进行Spark Jar作业开发前,请准备以下开发环境。表1 Spark Jar作业开发环境准备项说明操作系统Windows系统,支持Windows7以上版本。安装JDKJDK使用1.8版本。安装和配置IntelliJ IDEAIntelliJ IDEA为进行应用开发的工具,版本要求使用2019.1或其他兼容版本。安装Maven开发环境的基本配置。用于项目管理,贯穿软件开发生命周期。开发流程DLI进行Spark Jar作业开发流程参考如下:图1 Spark Jar作业开发流程表2 开发流程说明序号阶段操作界面说明1创建DLI通用队列DLI控制台创建作业运行的DLI队列。2上传数据到OBS桶OBS控制台将测试数据上传到OBS桶下。3新建Maven工程,配置pom文件IntelliJ IDEA参考样例代码说明,编写程序代码读取OBS数据。4编写程序代码5调试,编译代码并导出Jar包6上传Jar包到OBS和DLIOBS控制台将生成的Spark Jar包文件上传到OBS目录下和DLI程序包中。7创建Spark Jar作业DLI控制台在DLI控制台创建Spark Jar作业并提交运行作业。8查看作业运行结果DLI控制台查看作业运行状态和作业运行日志。步骤1:创建DLI通用队列第一次提交Spark作业,需要先创建队列,例如创建名为“sparktest”的队列,队列类型选择为“通用队列”。在DLI管理控制台的左侧导航栏中,选择“队列管理”。单击“队列管理”页面右上角“购买队列”进行创建队列。创建名为“sparktest”的队列,队列类型选择为“通用队列”。创建队列详细介绍请参考创建队列。单击“立即购买”,确认配置。配置确认无误,单击“提交”完成队列创建。步骤2:上传数据到OBS桶根据如下数据,创建people.json文件。{"name":"Michael"} {"name":"Andy", "age":30} {"name":"Justin", "age":19}进入OBS管理控制台,在“桶列表”下,单击已创建的OBS桶名称,本示例桶名为“dli-test-obs01”,进入“概览”页面。单击左侧列表中的“对象”,选择“上传对象”,将people.json文件上传到OBS桶根目录下。在OBS桶根目录下,单击“新建文件夹”,创建名为“result”的文件夹。单击“result”的文件夹,在“result”下单击“新建文件夹”,创建名为“parquet”的文件夹。步骤3:新建Maven工程,配置pom依赖以下通过IntelliJ IDEA 2020.2工具操作演示。打开IntelliJ IDEA,选择“File > New > Project”。图2 新建Project选择Maven,Project SDK选择1.8,单击“Next”。定义样例工程名和配置样例工程存储路径,单击“Finish”完成工程创建。如上图所示,本示例创建Maven工程名为:SparkJarObs,Maven工程路径为:“D:\DLITest\SparkJarObs”。在pom.xml文件中添加如下配置。<dependencies> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.11</artifactId> <version>2.3.2</version> </dependency> </dependencies>图3 修改pom.xml文件在工程路径的“src > main > java”文件夹上鼠标右键,选择“New > Package”,新建Package和类文件。Package根据需要定义,本示例定义为:“com.huawei.dli.demo”,完成后回车。在包路径下新建Java Class文件,本示例定义为:SparkDemoObs。步骤4:编写代码编写SparkDemoObs程序读取OBS桶下的1的“people.json”文件,并创建和查询临时表“people”。完整的样例请参考完整样例代码参考,样例代码分段说明如下:导入依赖的包。import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SaveMode; import org.apache.spark.sql.SparkSession; import static org.apache.spark.sql.functions.col;通过当前帐号的AK和SK创建SparkSession会话spark 。SparkSession spark = SparkSession .builder() .config("spark.hadoop.fs.obs.access.key", "xxx") .config("spark.hadoop.fs.obs.secret.key", "yyy") .appName("java_spark_demo") .getOrCreate();"spark.hadoop.fs.obs.access.key"参数对应的值"xxx"需要替换为帐号的AK值。"spark.hadoop.fs.obs.secret.key"参数对应的值“yyy”需要替换为帐号的SK值。AK和SK值获取请参考:如何获取AK和SK。读取OBS桶中的“people.json”文件数据。其中“dli-test-obs01”为演示的OBS桶名,请根据实际的OBS桶名替换。Dataset<Row> df = spark.read().json("obs://dli-test-obs01/people.json"); df.printSchema();通过创建临时表“people”读取文件数据。df.createOrReplaceTempView("people");查询表“people”数据。Dataset<Row> sqlDF = spark.sql("SELECT * FROM people"); sqlDF.show();将表“people”数据以parquet格式输出到OBS桶的“result/parquet”目录下。sqlDF.write().mode(SaveMode.Overwrite).parquet("obs://dli-test-obs01/result/parquet"); spark.read().parquet("obs://dli-test-obs01/result/parquet").show();关闭SparkSession会话spark。spark.stop();步骤5:调试、编译代码并导出Jar包单击IntelliJ IDEA工具右侧的“Maven”,参考下图分别单击“clean”、“compile”对代码进行编译。编译成功后,单击“package”对代码进行打包。打包成功后,生成的Jar包会放到target目录下,以备后用。本示例将会生成到:“D:\DLITest\SparkJarObs\target”下名为“SparkJarObs-1.0-SNAPSHOT.jar”。步骤6:上传Jar包到OBS和DLI下登录OBS控制台,将生成的“SparkJarObs-1.0-SNAPSHOT.jar”Jar包文件上传到OBS路径下。将Jar包文件上传到DLI的程序包管理中,方便后续统一管理。登录DLI管理控制台,单击“数据管理 > 程序包管理”。在“程序包管理”页面,单击右上角的“创建”创建程序包。在“创建程序包”对话框,配置以下参数。包类型:选择“JAR”。OBS路径:程序包所在的OBS路径。分组设置和组名称根据情况选择设置,方便后续识别和管理程序包。单击“确定”,完成创建程序包。步骤7:创建Spark Jar作业登录DLI控制台,单击“作业管理 > Spark作业”。在“Spark作业”管理界面,单击“创建作业”。在作业创建界面,配置对应作业运行参数。具体说明如下:所属队列:选择已创建的DLI通用队列。例如当前选择步骤1:创建DLI通用队列创建的通用队列“sparktest”。作业名称(--name):自定义Spark Jar作业运行的名称。当前定义为:SparkTestObs。应用程序:选择步骤6:上传Jar包到OBS和DLI下中上传到DLI程序包。例如当前选择为:“SparkJarObs-1.0-SNAPSHOT.jar”。主类:格式为:程序包名+类名。例如当前为:com.huawei.dli.demo.SparkDemoObs。其他参数可暂不选择,想了解更多Spark Jar作业提交说明可以参考创建Spark作业。图4 创建Spark Jar作业单击“执行”,提交该Spark Jar作业。在Spark作业管理界面显示已提交的作业运行状态。步骤8:查看作业运行结果在Spark作业管理界面显示已提交的作业运行状态。初始状态显示为“启动中”。如果作业运行成功则作业状态显示为“已成功”,单击“操作”列“更多”下的“Driver日志”,显示当前作业运行的日志。图5 “Driver日志”中的作业执行日志如果作业运行成功,本示例进入OBS桶下的“result/parquet”目录,查看已生成预期的parquet文件。如果作业运行失败,单击“操作”列“更多”下的“Driver日志”,显示具体的报错日志信息,根据报错信息定位问题原因。例如,如下截图信息因为创建Spark Jar作业时主类名没有包含包路径,报找不到类名“SparkDemoObs”。可以在“操作”列,单击“编辑”,修改“主类”参数为正确的:com.huawei.dli.demo.SparkDemoObs,单击“执行”重新运行该作业即可。后续指引如果您想通过Spark Jar作业访问其他数据源,请参考《使用Spark作业跨源访问数据源》。如果您想通过Spark Jar作业在DLI创建数据库和表,请参考《使用Spark作业访问DLI元数据》。完整样例代码参考package com.huawei.dli.demo; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SaveMode; import org.apache.spark.sql.SparkSession; import static org.apache.spark.sql.functions.col; public class SparkDemoObs { public static void main(String[] args) { SparkSession spark = SparkSession .builder() .config("spark.hadoop.fs.obs.access.key", "xxx") .config("spark.hadoop.fs.obs.secret.key", "yyy") .appName("java_spark_demo") .getOrCreate(); // can also be used --conf to set the ak sk when submit the app // test json data: // {"name":"Michael"} // {"name":"Andy", "age":30} // {"name":"Justin", "age":19} Dataset<Row> df = spark.read().json("obs://dli-test-obs01/people.json"); df.printSchema(); // root // |-- age: long (nullable = true) // |-- name: string (nullable = true) // Displays the content of the DataFrame to stdout df.show(); // +----+-------+ // | age| name| // +----+-------+ // |null|Michael| // | 30| Andy| // | 19| Justin| // +----+-------+ // Select only the "name" column df.select("name").show(); // +-------+ // | name| // +-------+ // |Michael| // | Andy| // | Justin| // +-------+ // Select people older than 21 df.filter(col("age").gt(21)).show(); // +---+----+ // |age|name| // +---+----+ // | 30|Andy| // +---+----+ // Count people by age df.groupBy("age").count().show(); // +----+-----+ // | age|count| // +----+-----+ // | 19| 1| // |null| 1| // | 30| 1| // +----+-----+ // Register the DataFrame as a SQL temporary view df.createOrReplaceTempView("people"); Dataset<Row> sqlDF = spark.sql("SELECT * FROM people"); sqlDF.show(); // +----+-------+ // | age| name| // +----+-------+ // |null|Michael| // | 30| Andy| // | 19| Justin| // +----+-------+ sqlDF.write().mode(SaveMode.Overwrite).parquet("obs://dli-test-obs01/result/parquet"); spark.read().parquet("obs://dli-test-obs01/result/parquet").show(); spark.stop(); } }
总条数:1488 到第
上滑加载中