-
MRS 3.1.0普通集群怎么安装开源的flink组件,求助各位大神
-
+mrs的端口该写什么mrs端口太多
-
前言 很多大数据计算都是用SQL实现的,跑得慢时就要去优化SQL,但常常碰到让人干瞪眼的情况。 很多大数据计算都是用SQL实现的,跑得慢时就要去优化SQL,但常常碰到让人干瞪眼的情况。 比如,存储过程中有三条大概形如这样的语句执行得很慢: select a,b,sum(x) from T group by a,b where …; select c,d,max(y) from T group by c,d where …; select a,c,avg(y),min(z) from T group by a,c where …; 这里的T是个有数亿行的巨大表,要分别按三种方式分组,分组的结果集都不大。 分组运算要遍历数据表,这三句SQL就要把这个大表遍历三次,对数亿行数据遍历一次的时间就不短,何况三遍。 这种分组运算中,相对于遍历硬盘的时间,CPU计算时间几乎可以忽略。如果可以在一次遍历中把多种分组汇总都计算出来,虽然CPU计算量并没有变少,但能大幅减少硬盘读取数据量,就能成倍提速了。 如果SQL支持类似这样的语法: from T --数据来自T表 select a,b,sum(x) group by a,b where … --遍历中的第一种分组 select c,d,max(y) group by c,d where … --遍历中的第二种分组 select a,c,avg(y),min(z) group by a,c where …; --遍历中的第三种分组 能一次返回多个结果集,那就可以大幅提高性能了。 可惜, SQL没有这种语法,写不出这样的语句,只能用个变通的办法,就是用group a,b,c,d的写法先算出更细致的分组结果集,但要先存成一个临时表,才能进一步用SQL计算出目标结果。SQL大致如下: create table T\_temp as select a,b,c,d, sum(case when … then x else 0 end) sumx, max(case when … then y else null end) maxy, sum(case when … then y else 0 end) sumy, count(case when … then 1 else null end) county, min(case when … then z else null end) minz group by a,b,c,d; select a,b,sum(sumx) from T\_temp group by a,b where …; select c,d,max(maxy) from T\_temp group by c,d where …; select a,c,sum(sumy)/sum(county),min(minz) from T\_temp group by a,c where …; 这样只要遍历一次了,但要把不同的WHERE条件转到前面的case when里,代码复杂很多,也会加大计算量。而且,计算临时表时分组字段的个数变得很多,结果集就有可能很大,最后还对这个临时表做多次遍历,计算性能也快不了。大结果集分组计算还要硬盘缓存,本身性能也很差。 还可以用存储过程的数据库游标把数据一条一条fetch出来计算,但这要全自己实现一遍WHERE和GROUP的动作了,写起来太繁琐不说,数据库游标遍历数据的性能只会更差! 只能干瞪眼! TopN运算同样会遇到这种无奈。举个例子,用Oracle的SQL写top5大致是这样的: select \* from (select x from T order by x desc) where rownum<=5 1 表T有10亿条数据,从SQL语句来看,是将全部数据大排序后取出前5名,剩下的排序结果就没用了!大排序成本很高,数据量很大内存装不下,会出现多次硬盘数据倒换,计算性能会非常差! 避免大排序并不难,在内存中保持一个5条记录的小集合,遍历数据时,将已经计算过的数据前5名保存在这个小集合中,取到的新数据如果比当前的第5名大,则插入进去并丢掉现在的第5名,如果比当前的第5名要小,则不做动作。这样做,只要对10亿条数据遍历一次即可,而且内存占用很小,运算性能会大幅提升。 这种算法本质上是把TopN也看作与求和、计数一样的聚合运算了,只不过返回的是集合而不是单值。SQL要是能写成这样,就能避免大排序了: select top(x,5) from T 1 然而非常遗憾,SQL没有显式的集合数据类型,聚合函数只能返回单值,写不出这种语句! 不过好在全集的TopN比较简单,虽然SQL写成那样,数据库却通常会在工程上做优化,采用上述方法而避免大排序。所以Oracle算那条SQL并不慢。 但是,如果TopN的情况复杂了,用到子查询中或者和JOIN混到一起的时候,优化引擎通常就不管用了。比如要在分组后计算每组的TopN,用SQL写出来都有点困难。Oracle的SQL写出来是这样: select \* from (select y,x,row\_number() over (partition by y order by x desc) rn from T) where rn<=5 1 这时候,数据库的优化引擎就晕了,不会再采用上面说的把TopN理解成聚合运算的办法。只能去做排序了,结果运算速度陡降! 假如SQL的分组TopN能这样写: select y,top(x,5) from T group by y 1 把top看成和sum一样的聚合函数,这不仅更易读,而且也很容易高速运算。 可惜,不行。 还是干瞪眼! 关联计算也是很常见的情况。以订单和多个表关联后做过滤计算为例,SQL大体是这个样子: select o.oid,o.orderdate,o.amount from orders o left join city ci on o.cityid = ci.cityid left join shipper sh on o.shid=sh.shid left join employee e on o.eid=e.eid left join supplier su on o.suid=su.suid where ci.state='New York' and e.title='manager' and ... 订单表有几千万数据,城市、运货商、雇员、供应商等表数据量都不大。过滤条件字段可能会来自于这些表,而且是前端传参数到后台的,会动态变化。 SQL一般采用HASH JOIN算法实现这些关联,要计算 HASH 值并做比较。每次只能解析一个JOIN,有N个JOIN要执行N遍动作,每次关联后都需要保持中间结果供下一轮使用,计算过程复杂,数据也会被遍历多次,计算性能不好。 通常,这些关联的代码表都很小,可以先读入内存。如果将订单表中的各个关联字段预先做序号化处理,比如将雇员编号字段值转换为对应雇员表记录的序号。那么计算时,就可以用雇员编号字段值(也就是雇员表序号),直接取内存中雇员表对应位置的记录,性能比HASH JOIN快很多,而且只需将订单表遍历一次即可,速度提升会非常明显! 也就是能把SQL写成下面的样子: select o.oid,o.orderdate,o.amount from orders o left join city c on o.cid = c.# --订单表的城市编号通过序号#关联城市表 left join shipper sh on o.shid=sh.# --订单表运货商号通过序号#关联运货商表 left join employee e on o.eid=e.# --订单表的雇员编号通过序号#关联雇员表 left join supplier su on o.suid=su.#--订单表供应商号通过序号#关联供应商表 where ci.state='New York' and e.title='manager' and ... 1 2 3 4 5 6 7 8 9 可惜的是,SQL 使用了无序集合概念,即使这些编号已经序号化了,数据库也无法利用这个特点,不能在对应的关联表这些无序集合上使用序号快速定位的机制,只能使用索引查找,而且数据库并不知道编号被序号化了,仍然会去计算 HASH 值和比对,性能还是很差! 有好办法也实施不了,只能再次干瞪眼! 还有高并发帐户查询,这个运算倒是很简单: select id,amt,tdate,… from T where id='10100' and tdate>= to\_date('2021-01-10','yyyy-MM-dd') and tdate and="" …="" 1 2 3 4 5 在T表的几亿条历史数据中,快速找到某个帐户的几条到几千条明细,SQL写出来并不复杂,难点是大并发时响应速度要达到秒级甚至更快。为了提高查询响应速度,一般都会对 T 表的 id 字段建索引: create index index_T_1 on T(id) 1 在数据库中,用索引查找单个帐户的速度很快,但并发很多时就会明显变慢。原因还是上面提到的SQL无序理论基础,总数据量很大,无法全读入内存,而数据库不能保证同一帐户的数据在物理上是连续存放的。硬盘有最小读取单位,在读不连续数据时,会取出很多无关内容,查询就会变慢。高并发访问的每个查询都慢一点,总体性能就会很差了。在非常重视体验的当下,谁敢让用户等待十秒以上?! 容易想到的办法是,把几亿数据预先按照帐户排序,保证同一帐户的数据连续存储,查询时从硬盘上读出的数据块几乎都是目标值,性能就会得到大幅提升。 但是,采用SQL体系的关系数据库并没有这个意识,不会强制保证数据存储的物理次序!这个问题不是SQL语法造成的,但也和SQL的理论基础相关,在关系数据库中还是没法实现这些算法。 那咋办?只能干瞪眼吗? 不能再用SQL和关系数据库了,要使用别的计算引擎。 开源的集算器SPL基于创新的理论基础,支持更多的数据类型和运算,能够描述上述场景中的新算法。用简单便捷的SPL写代码,在短时间内能大幅提高计算性能! 上面这些问题用SPL写出来的代码样例如下: 一次遍历计算多种分组 A B 1 =file(“T.ctx”).open().cursor(a,b,c,d,x,y,z 2 cursor A1 =A2.select(…).groups(a,b;sum(x)) 3 //定义遍历中的第一种过滤、分组 4 cursor =A4.select(…).groups(c,d;max(y)) 5 //定义遍历中的第二种过滤、分组 6 cursor =A6.select(…).groupx(a,c;avg(y),min(z)) 7 //定义遍历中的第三种过滤、分组 8 … //定义结束,开始计算三种方式的过滤、分组 用聚合的方式计算Top5 全集Top5(多线程并行计算) A 1 =file(“T.ctx”).open() 2 =A1.cursor@m(x).total(top(-5,x),top(5,x)) 3 //top(-5,x) 计算出 x 最大的前 5 名,top(5,x) 是 x 最小的前 5 名。 分组Top5(多线程并行计算) A 1 =file(“T.ctx”).open() 2 =A1.cursor@m(x,y).groups(y;top(-5,x),top(5,x)) 用序号做关联的SPL代码: 系统初始化 A 1 >env(city,file(“city.btx”).import@b()),env(employee,file(“employee.btx”).import@b()),… 2 //系统初始化时,几个小表读入内存 查询 A 1 =file(“orders.ctx”).open().cursor(cid,eid,…).switch(cid,city:#;eid,employee:#;…) 2 =A1.select(cid.state==“New York” && eid.title==“manager”…) 3 //先序号关联,再引用关联表字段写过滤条件 高并发帐户查询的SPL代码: 数据预处理,有序存储 A B 1 =file(“T-original.ctx”).open().cursor(id,tdate,amt,…) 2 =A1.sortx(id) =file(“T.ctx”) 3 =B2.create@r(#id,tdate,amt,…).append@i(A2) 4 =B2.open().index(index_id;id) 5 //将原数据排序后,另存为新表,并为帐号建立索引 帐户查询 A 1 =T.icursor(;id==10100 && tdate>=date(“2021-01-10”) && tdate2 //查询代码非常简单 除了这些简单例子,SPL还能实现更多高性能算法,比如有序归并实现订单和明细之间的关联、预关联技术实现多维分析中的多层维表关联、位存储技术实现上千个标签统计、布尔集合技术实现多个枚举值过滤条件的查询提速、时序分组技术实现复杂的漏斗分析等等。 正在为SQL性能优化头疼的小伙伴们,可以和我们一起探讨: http://www.raqsoft.com.cn/wx/Query-run-batch-ad.html SPL资料 SPL官网 SPL下载 SPL源代码 ———————————————— 版权声明:本文为CSDN博主「IT邦德」的原创文章,遵循CC 4.0 BY-SA版权协议,转载请附上原文出处链接及本声明。 原文链接:https://blog.csdn.net/weixin_41645135/article/details/127099888
-
实验目的1. 理解分类算法的概念以及基本思想流程;2. 熟练使用Eclipse等环境编写调试分类算法的Java代码。 实验平台操作系统:WindowsEclipse或MyEclipse 实验内容3.1 分类算法思想与流程(自己截图或文字说明)3.2 Eclipse或MyEclipse编写调试分类算法的Java代码,分析运行结果。算法思想及过程:我们统计了14天的气象数据(指标包括outlook,temperature,humidity,windy),并已知这些天气是否打球(play)。如果给出新一天的气象指标数据:sunny,cool,high,TRUE,判断一下会不会去打球。table 1outlooktemperaturehumiditywindyplaysunnyhothighFALSEnosunnyhothighTRUEnoovercasthothighFALSEyesrainymildhighFALSEyesrainycoolnormalFALSEyesrainycoolnormalTRUEnoovercastcoolnormalTRUEyessunnymildhighFALSEnosunnycoolnormalFALSEyesrainymildnormalFALSEyessunnymildnormalTRUEyesovercastmildhighTRUEyesovercasthotnormalFALSEyesrainymildhighTRUEno这个问题当然可以用朴素贝叶斯法求解,分别计算在给定天气条件下打球和不打球的概率,选概率大者作为推测结果。现在我们使用ID3归纳决策树的方法来求解该问题。预备知识:信息熵熵是无序性(或不确定性)的度量指标。假如事件A的全概率划分是(A1,A2,...,An),每部分发生的概率是(p1,p2,...,pn),那信息熵定义为:通常以2为底数,所以信息熵的单位是bit。补充两个对数去处公式:ID3算法构造树的基本想法是随着树深度的增加,节点的熵迅速地降低。熵降低的速度越快越好,这样我们有望得到一棵高度最矮的决策树。在没有给定任何天气信息时,根据历史数据,我们只知道新的一天打球的概率是9/14,不打的概率是5/14。此时的熵为:属性有4个:outlook,temperature,humidity,windy。我们首先要决定哪个属性作树的根节点。对每项指标分别统计:在不同的取值下打球和不打球的次数。table 2outlooktemperaturehumiditywindyplay yesno yesno yesno yesnoyesnosunny23hot22high34FALSE6295overcast40mild42normal61TRUR33 rainy32cool31 下面我们计算当已知变量outlook的值时,信息熵为多少。outlook=sunny时,2/5的概率打球,3/5的概率不打球。entropy=0.971outlook=overcast时,entropy=0outlook=rainy时,entropy=0.971而根据历史统计数据,outlook取值为sunny、overcast、rainy的概率分别是5/14、4/14、5/14,所以当已知变量outlook的值时,信息熵为:5/14 × 0.971 + 4/14 × 0 + 5/14 × 0.971 = 0.693这样的话系统熵就从0.940下降到了0.693,信息增溢gain(outlook)为0.940-0.693=0.247同样可以计算出gain(temperature)=0.029,gain(humidity)=0.152,gain(windy)=0.048。gain(outlook)最大(即outlook在第一步使系统的信息熵下降得最快),所以决策树的根节点就取outlook。接下来要确定N1取temperature、humidity还是windy?在已知outlook=sunny的情况,根据历史数据,我们作出类似table 2的一张表,分别计算gain(temperature)、gain(humidity)和gain(windy),选最大者为N1。依此类推,构造决策树。当系统的信息熵降为0时,就没有必要再往下构造决策树了,此时叶子节点都是纯的--这是理想情况。最坏的情况下,决策树的高度为属性(决策变量)的个数,叶子节点不纯(这意味着我们要以一定的概率来作出决策)。 3 ID3算法Java实现 ID3算法实现包括四个类的设计: 一、 决策树节点类(TreeNode类),包括类属性:name(节点属性名称),rule(节点属性值域,也就是对应决策树的分裂规则),child(节点下的孩子节点),datas(当前决策下对应的样本元组), candidateAttr(当前决策下剩余的分类属性)。 二、 最大信息增益节点计算类(Gain类):包括属性值:D(当前决策层次下的样本数据),attrList(当前决策层次下的剩余分类属性);包括方法:统计属性取值方法,统计属性不同取值计数方法,计算先验熵和条件熵的方法,筛选指定属性索引在指定值上的样本元组方法,通过先验熵减后验熵计算出最大信息增益值属性的方法。具体方法在程序中都已经注释,在这里只是根据需求给出方法的大致功能。 三、决策树建立类(DecisionTree类):包括方法:计算当前样本中分类属性的取值及其计数,并由此计算出多数类,决策树节点递归构建构成,具体实现思想同课上讲授内容,在此不在重述,借助的类是增益值计算类。 四、 ID3算法测试类(TestDecisionTree类):借助于上面决策树建立类,决策树节点之间连接已经建立完毕,下面将以上第二部分的样本数据作为测试数据,并且实现递归打印方法,输出决策树具体内容。Java代码: import java.util.List; //计算经验熵 public class Entroy { public double HD(List list, Object[] category){ int n = list.get(0).length; double[] p = new double[category.length]; double HD = 0; for(int k = 0; k double num = 0; for(int i = 0; i < list.size(); i ++){ if(list.get(i)[n-1] == category[k]){ num++; } } p[k] = num/list.size(); } for(int m = 0; m < category.length; m ++){ HD = HD + (-1) * p[m] * (Math.log(p[m])/Math.log(2.0)); } return HD; } } import java.util.List; //计算条件熵 public class Conditition { public double GD(List list,List array, int n){ Object[] objects = array.get(n); double gd = 0; double[] p = new double[objects.length]; for(int k = 0; k double num = 0; for (int i = 0; i < list.size(); i++) { if (list.get(i)[n] == objects[k]) { num++; } } p[k] = num/list.size(); double h = HH(list,objects[k],n,array.get(array.size()-1)); gd = gd + p[k] * h; } return gd; } public double HH(List list, Object object, int n, Object[] category){ int m = list.get(0).length; double HD = 0; double[] p = new double[category.length]; for(int k = 0; k < category.length; k ++){ double num = 0, nums = 0; for(int i = 0; i < list.size(); i ++){ if(list.get(i)[n] == object){ num ++; if(list.get(i)[m-1] == category[k]){ nums ++; } } } p[k] = nums/num; } for(int j = 0; j < category.length; j ++) { if (p[j] == 0) { HD = 0; } else { HD = HD + (-1) * p[j] * (Math.log(p[j]) / Math.log(2.0)); } } return HD; } } //输出选择的特征,返回该特征维度 public class OutPut { public int output(double[] GD, String[] feature){ int max = 0; for(int i = 1; i < GD.length; i ++){ if(GD[i] > GD[max]){ max = i; } } System.out.println(feature[max]); return max; } } import java.util.ArrayList; import java.util.List; import static java.lang.Float.NaN; //ID3算法 import java.util.ArrayList;import java.util.List;import static java.lang.Float.NaN;public class ID3 { public void id3(List list, String[] feature , List array){ Entroy e = new Entroy(); int len = list.get(0).length; Conditition c = new Conditition(); Object[] category = array.get(len-1); double HD = e.HD(list,category); double[] GD = new double[array.size()-1]; double[] HDA = new double[array.size()-1]; for(int i =0; i HDA[i] = c.GD(list, array,i ); GD[i] = HD - HDA[i]; } OutPut outPut = new OutPut(); int max = outPut.output(GD,feature); List>lists = new ArrayList<>(); //将数据按照特征划分为不同区域,为下一步求熵值做准备 for(int k = 0; k < array.get(max).length; k ++) { List l = new ArrayList<>(); for (int i = 0; i < list.size(); i++) { if (list.get(i)[max] == array.get(max)[k]){ l.add(list.get(i)); } } lists.add(l); } boolean flag = false; for(int i = 0; i double[] GD1 = new double[array.size() - 1]; double[] HDA1 = new double[array.size() - 1]; for(int k = 0; k < lists.get(i).size(); k++){ if(lists.get(i).get(k)[len-1] != lists.get(i).get(0)[len-1]) { flag = true; } } if(flag){ id3(lists.get(i),feature,array); } } }}//测试import java.util.ArrayList;import java.util.List;import static java.lang.Float.NaN;public class test { public static void main(String[] args) { String[] feature = {"年龄","有工作","有自己的房子","信贷情况","类别"}; Object[] age ={"青年","中年","老年"}; Object[] work = {'是','否'}; Object[] house = {'是','否'}; Object[] loan = {"一般",'好',"非常好"}; Object[] category = {'是','否'}; List array = new ArrayList(); array.add(age); array.add(work); array.add(house); array.add(loan); array.add(category); Object[] o1 = {age[0],work[1],house[1],loan[0],category[1]}; Object[] o2 = {age[0],work[1],house[1],loan[1],category[1]}; Object[] o3 = {age[0],work[0],house[1],loan[1],category[0]}; Object[] o4 = {age[0],work[0],house[0],loan[0],category[0]}; Object[] o5 = {age[0],work[1],house[1],loan[0],category[1]}; Object[] o6 = {age[1],work[1],house[1],loan[0],category[1]}; Object[] o7 = {age[1],work[1],house[1],loan[1],category[1]}; Object[] o8 = {age[1],work[0],house[0],loan[1],category[0]}; Object[] o9 = {age[1],work[1],house[0],loan[2],category[0]}; Object[] o10 = {age[1],work[1],house[0],loan[2],category[0]}; Object[] o11 = {age[2],work[1],house[0],loan[2],category[0]}; Object[] o12 = {age[2],work[1],house[0],loan[1],category[0]}; Object[] o13 = {age[2],work[0],house[1],loan[1],category[0]}; Object[] o14 = {age[2],work[0],house[1],loan[2],category[0]}; Object[] o15 = {age[2],work[1],house[1],loan[0],category[1]}; List list = new ArrayList(); list.add(o1); list.add(o2); list.add(o3); list.add(o4); list.add(o5); list.add(o6); list.add(o7); list.add(o8); list.add(o9); list.add(o10); list.add(o11); list.add(o12); list.add(o13); list.add(o14); list.add(o15); ID3 id3 = new ID3(); id3.id3(list,feature,array); }}决策树算法-例2决策树:是一种用于对实例进行分类的树形结构,可以是二叉树或非二叉树,由节点(node)和有向边(directed edge)组成。其中每个非叶子节点表示一个特征属性,叶子节点代表类别属性,它的值由根节点到叶子节点这一分支的属性值确定。使用决策树进行分类的过程,就是从根节点出发,训练数据的分支走向,直到得到叶子节点的值停止计算,这时即可输出类别。 决策树算法是从数据的属性(或者特征)出发,以属性作为基础,划分不同的类。实现决策树的算法有很多种,有ID3、C4.5和CART等算法。下面我们介绍ID3算法。 二、ID3算法 ID3算法是由Quinlan首先提出的,该算法是以信息论为基础,以信息熵和信息增益为衡量标准,从而实现对数据的归纳分类。 算法原理:设 为训练样本,对于每个样本有m个属性,用A,B,C,……来表示每一个属性。样本类别为 算法计算步骤如下: 首先,构造分类决策树。 (1)计算训练样本D的信息熵,即 Pi 表示第i个类别个数占训练样本总数的比例。 (2)分别计算每个属性的条件熵,例如计算属性A相对于D的期望信息,即: 其中, 表示属性A将D划分为v个子集。 表示第j个子集的样本数比上样本总数。 (3)计算属性A的信息增益,即 同理,重复第二步、第三步,计算出其他属性的信息增益,直至所有属性计算完成。选择信息增益最大的属性为根节点,对样本集D进行第一次分裂;然后对余下的属性重复上述的步骤,选择作为第二层节点的属性,直到找到叶子节点为止。此时,就构造出一颗决策分类树。 决策树构造完成后,就可以对待分类的样本进行分类,得到类别。 三、ID3算法实例讲解 图中数据为训练样本,现在预测E= {天气=晴,温度=适中,湿度=正常,风速=弱} 的情况下活动是取消还是进行。属性为天气、温度、湿度、风速,类别为取消和进行。 样本D的信息熵, 在天气为晴时有5种情况,发现活动取消有3种,进行有2种,计算现在的条件熵: 天气为阴时有4种情况,活动进行的有4种,则条件熵为: 天气为雨时有5种情况,活动取消的有2种,进行的有3种,则条件熵为: 由于按照天气属性不同取值划分时,天气为晴占整个情况的5/14,天气为阴占整个情况的4/14,天气为雨占整个情况的5/14,则按照天气属性不同取值划分时的带权平均值熵为:算出的结果约为0.693. 则此时的信息增益Gain(活动,天气)= H(活动) - H(活动|天气) = 0.94- 0.693 = 0.246 同理我们可以计算出按照温度属性不同取值划分后的信息增益: Gain(活动,温度)= H(活动) - H(活动|温度) = 0.94- 0.911 = 0.029 按照湿度属性不同取值划分后的信息增益: Gain(活动,湿度)= H(活动) - H(活动|湿度) = 0.94- 0.789 = 0.151 按照风速属性不同取值划分后的信息增益: Gain(活动,风速)= H(活动) - H(活动|风速) = 0.94- 0.892 = 0.048 决策树的构造就是要选择当前信息增益最大的属性来作为当前决策树的节点。因此我们选择天气属性来做为决策树根节点,这时天气属性有3取值可能:晴,阴,雨,我们发现当天气为阴时,活动全为进行因此这件事情就可以确定了,而天气为晴或雨时,活动中有进行的也有取消的,事件还无法确定,这时就需要在剩下的属性中递归再次计算活动熵和信息增益,选择信息增益最大的属性来作为下一个节点,直到整个事件能够确定下来。 最后得到的决策树如下图所示: 所以,E= {天气=晴,温度=适中,湿度=正常,风速=弱} 的情况下活动是进行。 四、ID3算法Java实现 下面是实例的Java代码实现,算法实现前,需要转换为数字向量,其中,天气有晴,阴,雨,可分配2,1,0三个值,即晴=2,阴=1,雨=0;同理,风速,强= 1,弱=0;湿度,高=1,正常=0;温度,炎热=2,正常=1,寒冷=0。————————————————版权声明:本文为CSDN博主「XiaoXiao_Yang77」的原创文章,遵循CC 4.0 BY-SA版权协议,转载请附上原文出处链接及本声明。原文链接:https://blog.csdn.net/XiaoXiao_Yang77/article/details/79262704Java代码:import java.util.ArrayList;import java.util.Collections;import java.util.Comparator;import java.util.HashMap;import java.util.HashSet;import java.util.LinkedHashMap;import java.util.LinkedList;import java.util.List;import java.util.Map;import java.util.Map.Entry;import java.util.Set;/** * * @author X.H.Yang */public class TreeID3Act { public static String[] createDataLable() { String lable[] = {"weather", " temper", "humidity", "wind"}; return lable; } public static Object[][] createDataSet() { Object set[][] = {{2, 2, 1, 0,"no"}, {2, 2, 1, 1, "no"}, {1, 2, 1, 0, "yes"}, {0, 1, 1, 0,"yes"}, {0, 0, 0, 0, "yes"}, {0, 0, 0, 1, "no"}, {1, 0, 0, 1, "yes"}, {2, 1, 1, 0, "no"}, {2, 0, 0, 0, "yes"}, {0, 1, 0, 0, "yes"}, {2, 1, 0, 1, "yes"}, {1, 1, 1, 1, "yes"}, {1, 2, 0, 0, "yes"}, {0, 1, 1, 1, "no"}}; return set; } public static double calcShannonEnt(Object[][] dataSet) { double shannonEnt = 0.0; int numEntries = dataSet.length; HashMap labelCounts = new HashMap<>(); for (Object[] featVec : dataSet) { String currentLabel = (String) featVec[featVec.length - 1]; if (labelCounts.get(currentLabel) == null) { labelCounts.put(currentLabel, 1); } else { int i = labelCounts.get(currentLabel); i++; labelCounts.put(currentLabel, i); } } for (Entry entry : labelCounts.entrySet()) { double prob = (double) entry.getValue() / (double) numEntries; shannonEnt -= prob * (Math.log(prob) / Math.log(2.0d)); } return shannonEnt; } public static ArrayList splitDataSet(Object dataSet[][], int axis, Object value) { ArrayList retDataSet = new ArrayList<>(); for (Object[] featVec : dataSet) { Object subSet[] = null; if (featVec[axis].equals(value)) { subSet = new Object[dataSet[0].length - 1]; if (axis == 0) { System.arraycopy(featVec, 1, subSet, 0, subSet.length); } else { System.arraycopy(featVec, 0, subSet, 0, axis); System.arraycopy(featVec, axis + 1, subSet, axis, subSet.length - axis); } retDataSet.add(subSet); } } return retDataSet; } public static int chooseBestFeatureToSplit(Object dataSet[][]) { int numFeatures = dataSet[0].length - 1; double baseEntropy = calcShannonEnt(dataSet); double bestInfoGain = 0.0; int bestFeature = -1; for (int f = 0; f < numFeatures; f++) { HashSet uniqueVals = new HashSet<>(); for (Object set[] : dataSet) { uniqueVals.add(set[f]); } double newEntropy = 0.0; for (Object obj : uniqueVals) { ArrayList subDataSet = splitDataSet(dataSet, f, obj); double prob = (double) subDataSet.size() / (double) dataSet.length; Object subSetArray[][] = new Object[subDataSet.size()][numFeatures]; for (int i = 0; i < subSetArray.length; i++) { System.arraycopy(subDataSet.get(i), 0, subSetArray[i], 0, numFeatures); } newEntropy += prob * calcShannonEnt(subSetArray); } double infoGain = baseEntropy - newEntropy; if (infoGain > bestInfoGain) { bestInfoGain = infoGain; bestFeature = f; } } return bestFeature; } public static String majorityCnt(String classList[]) { HashMap classCount = new HashMap<>(); for (String vote : classList) { if (!classCount.containsKey(vote)) { classCount.put(vote, 1); } else { int i = classCount.get(vote); i++; classCount.put(vote, i); } } LinkedHashMap sortMap = sortMapByValues(classCount); return sortMap.entrySet().iterator().next().getKey(); } private static Object createTree(Object dataSet[][], String labels[]) { //classList = [example[-1] for example in dataSet] ArrayList classList = new ArrayList(); for (Object set[] : dataSet) { classList.add((String) set[set.length - 1]); } if (ListCount(classList, classList.get(0)) == classList.size()) { return classList.get(0); } if (dataSet[0].length == 1) { return majorityCnt((String[]) classList.toArray()); } int bestFeat = chooseBestFeatureToSplit(dataSet); String bestFeatLabel = labels[bestFeat]; HashMap myTree = new HashMap<>(); myTree.put(bestFeatLabel, new HashMap<>()); String sublabels[] = new String[labels.length - 1]; if (bestFeat == 0) { System.arraycopy(labels, 1, sublabels, 0, sublabels.length); } else { System.arraycopy(labels, 0, sublabels, 0, bestFeat); System.arraycopy(labels, bestFeat + 1, sublabels, bestFeat, sublabels.length - bestFeat); } HashSet uniqueVals = new HashSet<>(); for (Object set[] : dataSet) { uniqueVals.add(set[bestFeat]); } for(Object value : uniqueVals){ ArrayList setlist = splitDataSet(dataSet, bestFeat, value); int j = 0; Object dataSetM[][] =new Object[setlist.size()][]; for(Object set[]:setlist){ dataSetM[j] = set; j++; } Object tree = createTree(dataSetM,sublabels); HashMap subtree = (HashMap) myTree.get(bestFeatLabel); subtree.put(value, tree); } return myTree; } private static String classify(HashMap inputTree,String featLabels[],Object[] testVec){ String firstStr = (String) inputTree.keySet().iterator().next(); HashMap secondDict = (HashMap) inputTree.get(firstStr); int featIndex = 0; for(featIndex=0;featIndex if(featLabels[featIndex].equals(firstStr)){ break; } } Object key = testVec[featIndex]; Object valueOfFeat = secondDict.get(key); String classLabel ="erro"; if(valueOfFeat instanceof HashMap){ classLabel = classify((HashMap) valueOfFeat, featLabels, testVec); }else{ classLabel = (String)valueOfFeat; } return classLabel; } private static int ListCount(ArrayList classList, String key) { int c = 0; for (String clazz : classList) { if (clazz.equals(key)) { c++; } } return c; } private static LinkedHashMap sortMapByValues(Map aMap) { Set> mapEntries = aMap.entrySet(); //System.out.println("Values and Keys before sorting "); //for (Entry entry : mapEntries) { //System.out.println(entry.getValue() + " - " + entry.getKey()); //} // used linked list to sort, because insertion of elements in linked list is faster than an array list. List> aList = new LinkedList>(mapEntries); // sorting the List Collections.sort(aList, new Comparator>() { @Override public int compare(Entry ele1, Entry ele2) { return -(ele1.getValue().compareTo(ele2.getValue())); } }); // Storing the list into Linked HashMap to preserve the order of insertion. LinkedHashMap aMap2 = new LinkedHashMap(); for (Entry entry : aList) { aMap2.put(entry.getKey(), entry.getValue()); } return aMap2; } public static void main(String[] args) throws Exception { Object Dataset[][] = createDataSet(); Object tree = createTree(Dataset,new String[]{"weather", " temper", "humidity", "wind"}); System.out.println(tree); String clazz = classify((HashMap) tree,new String[]{"weather", " temper", "humidity", "wind"},new Object[]{2,1,0,0}); System.out.println(clazz); }}
-
我看了官方的demo代码,在hive to hbase项目代码里,只设置了appName,其余的全部没有设置,是可以自动读取hive-site.xml等配置文件吗?huaweicloud-mrs-example/SparkHivetoHbase.java at mrs-3.0.2 · huaweicloud/huaweicloud-mrs-example (github.com)这是我举例的代码连接这个是代码中读取hive表数据的代码片段 SparkConf conf = new SparkConf().setAppName("SparkHivetoHbase"); JavaSparkContext jsc = new JavaSparkContext(conf); HiveContext sqlContext = new org.apache.spark.sql.hive.HiveContext(jsc); Dataset dataFrame = sqlContext.sql("select name, account from person");如果在代码中需要设置的话我有一个问题,hive默认的元数据服务是DBService,那hive.metastore.uris这一项应该怎么配置
-
9月20日,在曼谷举行的华为全联接大会2022上,华为全面阐述了数据存储产业“以数据为中心”的创新理念,聚焦场景需求为行业场景找技术,推出场景化的存储产品和解决方案,助力企业释放数据生产力。华为认为,数字化时代,向数据要生产力,数据存储面临四大变化:从传统数据库到分布式数据库、大数据、人工智能等多样化新兴应用蓬勃发展。数据热度不断攀升,对数据分析处理的实时性要求越来越高。自然灾害、人为操作失误等频发,从防止物理性破坏到防止人因性破坏,提升企业数字化韧性是当务之急。数据存储的绿色节能成为常态化要求。华为数据存储产品线总裁周跃峰表示: 华为数据存储将积极拥抱变化,深入践行“以数据为中心”的理念,面向生产交易、数据分析、数据保护等数据应用场景,构建可靠、高效的存储底座。生产交易面向金融Core banking,医疗HIS等场景,OceanStor Dorado全闪存存储,在领先的SAN基础上全面增强了NAS能力,NAS和SAN共享FlashLink智能盘控配合算法和SmartMatrix 高可靠全互联架构。业界唯一的Active-Active NAS双活能力,为文件业务提供连续性保障。数据分析面向高性能数据分析HPDA,医疗PACS等场景,OceanStor Pacific分布式存储,通过大小IO自适应数据流、融合非结构化数据索引、超高密硬件和弹性EC算法等技术架构突破,打破数据分析的性能墙、协议墙和容量墙,大幅提升数据分析处理效率30%以上。数据保护对于自然灾害或者人为因素导致的破坏,华为存储提供全面的数据保护能力。在容灾上,提供本地高可用、同城双活、两地三中心等多种容灾方案。在备份上,OceanProtect专用备份存储实现业界3倍备份带宽、5倍恢复带宽、72:1数据缩减率。同时,华为提供从主存储到备份存储的全方位防勒索存储解决方案,使用机器学习 模型进行勒索软件检测,检测率达99%以上。自动化管理对于存储运维管理,传统方式耗时耗人,通过部署DME可以实现规-建-维-优全生命周期自动化运维管理,提前14天硬盘故障预测、提前90天容量预测,从而帮助运维人员从繁杂的例行工作中解放出来,投入到更加有创新性的工作中。容器存储容器是当前行业热门的创新技术,在互联网、金融等领域非常流行。OceanStor Dorado NAS 实现容器应用和存储解耦部署,应用和存储按需扩展;同时OceanStor Dorado NAS提供跨节点数据共享以及比业界标杆高30%的资源调度效率。多云随着企业云化演进的不断实践,多公有云加多私有云正逐渐成为向云演进的最佳选择。数据集中共享存储、应用部署在多云是企业面向多云演进的最佳IT架构,存储厂商将专业存储以软硬件一体或者纯软件的方式部署到公有云平台,帮助企业实现数据的跨云平滑演进。华为存储正在积极开展这方面的创新实践。周跃峰表示,我们正在迎来YB数据时代,数据应用蓬勃发展,华为数据存储将以数据为中心,为客户构建可靠存储底座,释放数据生产力,最大化数据价值。
-
数字经济时代的来临,推动市场向着更灵活高效的形态演进,也促进了一批新业态和新模式的形成。而数据作为这一时代的核心要素,凭借其强大的潜在生长力,改变了人们对于生产效力的认知,不断地为社会发展注入新的活力。 2020年,国家层面将“数据”作为新型生产要素写入中央文件中,“数据是资产”也已成为市场的共识。随着数据资产成为企业数字化转型中的重要基石,市场主体对于数据管理的需求也持续扩张。 而数据资产运营是挖掘数据价值的有效引擎,能够决定企业数字化转型蜕变之路的成败。由数据资源集成的大数据产业生态,需要高效的数据资产运营手段为伍,才能实现企业商业价值最大化。企业未来想实现数智化发展,就要充分认识数字时代的转型关键,重点把握核心数据资产运营环节,通过闭环管理、周期评估、优化迭代运营过程,最终形成动态可持续的数据应用价值链。 未来,数据资产运营将成为企业商业价值实现的重要一环,如何有效掌握这一关键,全面释放数据潜能,将成为各市场主体的战略突破点。 白皮书指出,数据资产运营的关键在于“三要素”与“四重奏”。 在数字经济时代背景下,企业通过数据资产运营实现数据价值的重要性凸显。数据资产如何运营?其核心在于向下扎根与向上生长。向下扎根,意味着围绕数据建立良好的组织与意识、流程与规范、平台与工具,以在公司内部培育数据资产运营土壤。向上生长,意味着在深深扎根数据文化土壤的基础上,通过开展数据资产的盘点、评估、治理与共享,将核心业务数据资产紧握手中,实现数据资产稳健运营。二者缺一不可,只有做到这两点,才能助力企业数据互联互通,释放独有价值。“三要素”之组织与意识:凝心聚力——建立完善组织架构体系,培养数据资产运营文化 组织与意识需要企业管理层的引领。构建基于数据驱动的组织模式,实现组织内外人、物、知识等资源弹性供给和单元的动态协作是构建数字企业的基础,也是组织适应、利用、驾驭不确定环境的利器,数据资产运营组织架构的调整帮助企业打造一体化柔性运营管控能力,针对公司需求作出更敏捷的响应。“三要素”之流程与规范:规圆矩方——数据管理规范化,确保规章制度有效实施,奠定创新数据管理基础 无规矩不成方圆,数据驱动型企业在创新拼搏的同时,也亟需对数据全生命周期进行有机管理。规范化的数据管理的流程与制度能够让数据的交互、整合与使用更加流畅,从而大大减少数据上的问题与冲突,助力公司的数据管理由传统模式健康稳定地转向创新模式。“三要素”之平台与工具:巧借东风——数据资产管理平台承载数据产业化与商品化 平台与工具意味着生产力,是开展数据资产管理不可或缺的底层基石。通过一体化的系统框架体系,集中治理数据问题、集中进行数据监控运维与服务运营,不仅将传统数据管理工具各个组件进行了整合,更是将其进行打通与融合,实现数据在平台上的有效运转。“四重奏”之数据资产盘点:细致入微——双视角厘清数据资产,绘制企业级数据资产地图 数据资产盘点是数据资产运营的先行任务,旨在解决数据资产“有什么”的问题。从业务视角与技术视角出发,形成企业数据资产框架和数据资产目录,支持建立全面覆盖的企业级数据资产地图,为数据资产“用什么”以及“如何用”奠定基础。“四重奏”之数据资产评估:详察形候——多维度评估企业数据资产 数据资产评估是时代赋予企业的课题。在现如今的数字经济时代,随着数据、算法的升级,“资产”的形态和范围正在出现革命性的变化,“数据资产”这一概念也应运而生。但目前这一资产形式尚未体现在企业的财务报表上,同时也面临着诸多争议与挑战,合理、健全、有效的数据资产评估方式对社会数字化趋势的意义尤为深远。“四重奏”之数据治理:准绳嘉量——全方位构建数据治理完整链条,奠定高质量数据基石 数据治理是对数据资产的管理形式权力和控制的活动集合,它不仅仅是一套用工具组合的产品级解决方案,更是从决策层到技术层,从管理制度到工具支撑,自上而下贯穿整个组织架构的完整链条,以期通过持续的评估、指导和监督,确保富有成效且高效的数据利用,促进组织协作和结构化决策,为企业创造价值。“四重奏”之数据共享:内外兼修——全面释放数据价值 数据共享是盘活企业数据资产的有效手段。企业通过开展数据资产的内部循环与外部流通,运用数据分析与挖掘获取新的信息,在企业内部形成数据流转与共享,在企业外部为社会提供数据资产的价值,也同时为企业谋取创新型的收益,实现数据的增值。 后疫情时代,企业数字化转型需求不断攀升,塑造线上化、数字化服务能力成为转型变革的核心诉求,而数据资产运营则是转型中的关键环节。 数生万物,转型之本。响应数字化趋势不能只靠口号,而是需要充分把握“三要素”与“四重奏”的关键要义,搭建规范完善的数据管理框架,借力数据资产运营持续激发数据活力。通过数据重塑业务链与管理链,用数据贯穿企业经营与创新发展,使数据成为真正的生产力。最终实现企业长远发展愿景。来源:中国大数据产业观察,如有侵权,请联系删除~
-
基于JVM的开源数据处理语言主要有Kotlin、Scala、SPL,下面对三者进行多方面的横向比较,从中找出开发效率最高的数据处理语言。本文的适用场景设定为项目开发中常见的数据处理和业务逻辑,以结构化数据为主,大数据和高性能不作为重点,也不涉及消息流、科学计算等特殊场景。 基本特征 适应面 Kotlin的设计初衷是开发效率更高的Java,可以适用于任何Java涉及的应用场景,除了常见的信息管理系统,还能用于WebServer、Android项目、游戏开发,通用性比较好。Scala的设计初衷是整合现代编程范式的通用开发语言,实践中主要用于后端大数据处理,其他类型的项目中很少出现,通用性不如Kotlin。SPL的设计初衷是专业的数据处理语言,实践与初衷一致,前后端的数据处理、大小数据处理都很适合,应用场景相对聚焦,通用性不如Kotlin。 编程范式 Kotlin以面向对象编程为主,也支持函数式编程。Scala两种范式都支持,面向对象编程比Koltin更彻底,函数式编程也比Koltin方便些。SPL可以说不算支持面向对象编程,有对象概念,但没有继承重载这些内容,函数式编程比Kotlin更方便。 运行模式 Kotlin和Scala是编译型语言,SPL是解释型语言。解释型语言更灵活,但相同代码性能会差一点。不过SPL有丰富且高效的库函数,总体性能并不弱,面对大数据时常常会更有优势。 外部类库 Kotlin可以使用所有的Java类库,但缺乏专业的数据处理类库。Scala也可以使用所有的Java类库,且内置专业的大数据处理类库(Spark)。SPL内置专业的数据处理函数,提供了大量时间复杂度更低的基本运算,通常不需要外部Java类库,特殊情况可在自定义函数中调用。 IDE和调试 三者都有图形化IDE和完整的调试功能。SPL的IDE专为数据处理而设计,结构化数据对象呈现为表格形式,观察更加方便,Kotlin和Scala的IDE是通用的,没有为数据处理做优化,无法方便地观察结构化数据对象。 学习难度 Kotlin的学习难度稍高于Java,精通Java者可轻易学会。Scala的目标是超越Java,学习难度远大于Java。SPL的目标就是简化Java甚至SQL的编码,刻意简化了许多概念,学习难度很低。 代码量 Kotlin的初衷是提高Java的开发效率,官方宣称综合代码量只有Java的20%,可能是数据处理类库不专业的缘故,这方面的实际代码量降低不多。Scala的语法糖不少,大数据处理类库比较专业,代码量反而比Kotlin低得多。SPL只用于数据处理,专业性最强,再加上解释型语言表达能力强的特点,完成同样任务的代码量远远低于前两者(后面会有对比例子),从另一个侧面也能说明其学习难度更低。 语法 数据类型 原子数据类型:三者都支持,比如Short、Int、Long、Float、Double、Boolean 日期时间类型:Kotlin缺乏易用的日期时间类型,一般用Java的。Scala和SPL都有专业且方便的日期时间类型。 有特色的数据类型:Kotlin支持非数值的字符Char、可空类型Any?。Scala支持元组(固定长度的泛型集合)、内置BigDecimal。SPL支持高性能多层序号键,内置BigDecimal。 集合类型:Kotlin和Scala支持Set、List、Map。SPL支持序列(有序泛型集合,类似List)。 结构化数据类型:Kotlin有记录集合List,但缺乏元数据,不够专业。Scala有专业的结构化数类型,包括Row、RDD、DataSet、DataFrame(本文以此为例进行说明)等。SPL有专业的结构化数据类型,包括record、序表(本文以此为例进行说明)、内表压缩表、外存Lazy游标等。 Scala独有隐式转换能力,理论上可以在任意数据类型之间进行转换(包括参数、变量、函数、类),可以方便地改变或增强原有功能。 流程处理 三者都支持基础的顺序执行、判断分支、循环,理论上可进行任意复杂的流程处理,这方面不多讨论,下面重点比较针对集合数据的循环结构是否方便。以计算比上期为例,Kotlin代码: mData.forEachIndexed{index,it-> if(index>0) it.Mom= it.Amount/mData[index-1].Amount-1 } Kotlin的forEachIndexed函数自带序号变量和成员变量,进行集合循环时比较方便,支持下标取记录,可以方便地进行跨行计算。Kotlin的缺点在于要额外处理数组越界。 Scala代码: val w = Window.orderBy(mData("SellerId")) mData.withColumn("Mom", mData ("Amount")/lag(mData ("Amount"),1).over(w)-1) Scala跨行计算不必处理数组越界,这一点比Kotlin方便。但Scala的结构化数据对象不支持下标取记录,只能用lag函数整体移行,这对结构化数据不够方便。lag函数不能用于通用性强的forEach,而要用withColumn之类功能单一的循环函数。为了保持函数式编程风格和SQL风格的底层统一,lag函数还必须配合窗口函数(Python的移行函数就没这种要求),整体代码看上去反而比Kotlin复杂。 SPL代码: mData.(Mom=Amount/Amount[-1]-1) SPL对结构化数据对象的流程控制进行了多项优化,类似forEach这种最通用最常用的循环函数,SPL可以直接用括号表达,简化到极致。SPL也有移行函数,但这里用的是更符合直觉的“[相对位置]"语法,进行跨行计算时比Kotlin的绝对定位强大,比Scala的移行函数方便。上述代码之外,SPL还有更多针对结构化数据的流程处理功能,比如:每轮循环取一批而不是一条记录;某字段值变化时循环一轮。 Lambda表达式 Lambda表达式是匿名函数的简单实现,目的是简化函数的定义,尤其是变化多样的集合计算类函数。Kotlin支持Lambda表达式,但因为编译型语言的关系,难以将参数表达式方便地指定为值参数或函数参数,只能设计复杂的接口规则进行区分,甚至有所谓高阶函数专用接口,这就导致Kotin的Lambda表达式编写困难,在数据处理方面专业性不足。几个例子: "abcd".substring( 1,2) //值参数 "abcd".sumBy{ it.toInt()} //函数参数 mData.forEachIndexed{ index,it-> if(index>0) it.Mom=…} //函数参数的函数带多个参数 Koltin的Lambda表达式专业性不足,还表现在使用字段时必须带上结构化数据对象的变量名(it),而不能像SQL那样单表计算时可以省略表名。 同为编译型语言,Scala的Lambda表达式和Kotlin区别不大,同样需要设计复杂的接口规则,同样编写困难,这里就不举例了。计算比上期时,字段前也要带上结构化数据对象变量名或用col函数,形如mData (“Amount”)或col(“Amount”),虽然可以用语法糖弥补,写成$”Amount”或’Amount,但很多函数不支持这种写法,硬要弥补反而使风格不统一。 SPL的Lambda表达式简单易用,比前两者更专业,这与其解释型语言的特性有关。解释型语言可以方便地推断出值参数和函数参数,没有所谓复杂的高阶函数专用接口,所有的函数接口都一样简单。几个例子: mid("abcd",2,1) //值参数 Orders.sum(Amount*Amount) //函数参数 mData.(Mom=Amount/Amount[-1]-1) //函数参数的函数带多个参数 1 2 3 SPL可直接使用字段名,无须结构化数据对象变量名,比如: Orders.select(Amount>1000 && Amount<=3000 && like(Client,"*S*")) 1 SPL的大多数循环函数都有默认的成员变量~和序号变量#,可以显著提升代码编写的便利性,特别适合结构化数据计算。比如,取出偶数位置的记录: Students.select(# % 2==0) 1 求各组的前3名: Orders.group(SellerId;~.top(3;Amount)) 1 SPL函数选项和层次参数 值得一提的是,为了进一步提高开发效率,SPL还提供了独特的函数语法。 有大量功能类似的函数时,大部分程序语言只能用不同的名字或者参数进行区分,使用不太方便。而SPL提供了非常独特的函数选项,使功能相似的函数可以共用一个函数名,只用函数选项区分差别。比如,select函数的基本功能是过滤,如果只过滤出符合条件的第1条记录,可使用选项@1: T.select@1(Amount>1000) 1 对有序数据用二分法进行快速过滤,使用@b: T.select@b(Amount>1000) 1 函数选项还可以组合搭配,比如: Orders.select@1b(Amount>1000) 1 有些函数的参数很复杂,可能会分成多层。常规程序语言对此并没有特别的语法方案,只能生成多层结构数据对象再传入,非常麻烦。SQL使用了关键字把参数分隔成多个组,更直观简单,但这会动用很多关键字,使语句结构不统一。而SPL创造性地发明了层次参数简化了复杂参数的表达,通过分号、逗号、冒号自高而低将参数分为三层: join(Orders:o,SellerId ; Employees:e,EId) 1 数据源 数据源种类 Kotlin原则上可以支持所有的Java数据源,但代码很繁琐,类型转换麻烦,稳定性也差,这是因为Kotlin没有内置的数据源访问接口,更没有针对结构化数据处理做优化(JDBC接口除外)。从这个意义讲,也可以说它不直接支持任何数据源,只能使用Java第三方类库,好在第三方类库的数量足够庞大。 Scala支持的数据源种类比较多,且有六种数据源接口是内置的,并针对结构化数据处理做了优化,包括:JDBC、CSV、TXT、JSON、Parquet列存格式、ORC列式存储,其他的数据源接口虽然没有内置,但可以用社区小组开发的第三方类库。Scala提供了数据源接口规范,要求第三方类库输出为结构化数据对象,常见的第三方接口有XML、Cassandra、HBase、MongoDB等。 SPL内置了最多的数据源接口,并针对结构化数据处理做了优化,包括: JDBC(即所有的RDB) CSV、TXT、JSON、XML、Excel HBase、HDFS、Hive、Spark Salesforce、阿里云 Restful、WebService、Webcrawl Elasticsearch、MongoDB、Kafka、R2dbc、FTP Cassandra、DynamoDB、influxDB、Redis、SAP 这些数据源都可以直接使用,非常方便。对于其他未列入的数据源,SPL也提供了接口规范,只要按规范输出为SPL的结构化数据对象,就可以进行后续计算。 代码比较 以规范的CSV文件为例,比较三种语言的解析代码。Kotlin: val file = File("D:\\data\\Orders.txt") data class Order(var OrderID: Int,var Client: String,var SellerId: Int, var Amount: Double, var OrderDate: Date) var sdf = SimpleDateFormat("yyyy-MM-dd") var Orders=file.readLines().drop(1).map{ var l=it.split("\t") var r=Order(l[0].toInt(),l[1],l[2].toInt(),l[3].toDouble(),sdf.parse(l[4])) r } var resutl=Orders.filter{ it.Amount>= 1000 && it.Amount < 3000} Koltin专业性不足,通常要硬写代码读取CSV,包括事先定义数据结构,在循环函数中手工解析数据类型,整体代码相当繁琐。也可以用OpenCSV等类库读取,数据类型虽然不用在代码中解析,但要在配置文件中定义,实现过程不见得简单。 Scala专业性强,内置解析CSV的接口,代码比Koltin简短得多: val spark = SparkSession.builder().master("local").getOrCreate() val Orders = spark.read.option("header", "true").option("sep","\t").option("inferSchema", "true").csv("D:/data/orders.csv").withColumn("OrderDate", col("OrderDate").cast(DateType)) Orders.filter("Amount>1000 and Amount<=3000") Scala在解析数据类型时麻烦些,其他方面没有明显缺点。 SPL更加专业,连解析带计算只要一行: T("D:/data/orders.csv").select(Amount>1000 && Amount<=3000) 1 跨源计算 JVM数据处理语言的开放性强,有足够的能力对不同的数据源进行关联、归并、集合运算,但数据处理专业性的差异,导致不同语言的方便程度区别较大。 Kotlin不够专业,不仅缺乏内置数据源接口,也缺乏跨源计算函数,只能硬写代码实现。假设已经从不同数据源获得了员工表和订单表,现在把两者关联起来: data class OrderNew(var OrderID:Int ,var Client:String, var SellerId:Employee ,var Amount:Double ,var OrderDate:Date ) val result = Orders.map { o->var emp=Employees.firstOrNull{ it.EId==o.SellerId } emp?.let{ OrderNew(o.OrderID,o.Client,emp,o.Amount,o.OrderDate) } } .filter {o->o!=null} 很容易看出Kotlin的缺点,代码只要一长,Lambda表达式就变得难以阅读,还不如普通代码好理解;关联后的数据结构需要事先定义,灵活性差,影响解题流畅性。 Scala比Kotlin专业,不仅内置了多种数据源接口,而且提供了跨源计算的函数。同样的计算,Scala代码简单多了: val join=Orders.join(Employees,Orders("SellerId")===Employees("EId"),"Inner") 1 可以看到,Scala不仅具备专用于结构化数据计算的对象和函数,而且可以很好地配合Lambda语言,代码更易理解,也不用事先定义数据结构。 SPL更加专业,结构化数据对象更专业,跨源计算函数更方便,代码更简短: join(Orders:o,SellerId;Employees:e,EId) 1 自有存储格式 反复使用的中间数据,通常会以某种格式存为本地文件,以此提高取数性能。Kotlin支持多种格式的文件,理论上能够进行中间数据的存储和再计算,但因为在数据处理方面不专业,基本的读写操作都要写大段代码,相当于并没有自有的存储格式。 Scala支持多种存储格式,其中parquet文件常用且易用。parquet是开源存储格式,支持列存,可存储大量数据,中间计算结果(DataFrame)可以和parquet文件方便地互转。遗憾的是,parquet的索引尚不成熟。 val df = spark.read.parquet("input.parquet") val result=df.groupBy(data("Dept"),data("Gender")).agg(sum("Amount"),count("*")) result.write.parquet("output.parquet") SPL支持btx和ctx两种私有二进制存储格式,btx是简单行存,ctx支持行存、列存、索引,可存储大量数据并进行高性能计算,中间计算结果(序表/游标)可以和这两种文件方便地互转。 A 1 =file("input.ctx").open() 2 =A1.cursor(Dept,Gender,Amount).groups(Dept,Gender;sum(Amount):amt,count(1):cnt) 3 =file("output.ctx").create(#Dept,#Gender,amt,cnt).append(A2.cursor()) 结构化数据计算 结构化数据对象 数据处理的核心是计算,尤其是结构化数据的计算。结构化数据对象的专业程度,深刻地决定了数据处理的方便程度。 Kotlin没有专业的结构化数据对象,常用于结构化数据计算的是List,其中EntityBean可以用data class简化定义过程。 List是有序集合(可重复),凡涉及成员序号和集合的功能,Kotlin支持得都不错。比如按序号访问成员: Orders[3] //按下标取记录,从0开始 Orders.take(3) //前3条记录 Orders.slice(listOf(1,3,5)+IntRange(7,10)) //下标是1、3、5、7-10的记录 还可以按倒数序号取成员: Orders.reversed().slice(1,3,5) //倒数第1、3、5条 Orders.take(1)+Orders.takeLast(1) //第1条和最后1条 涉及顺序的计算难度都比较大,Kotlin支持有序计集合,进行相关的计算会比较方便。作为集合的一种,List擅长的功能还有集合成员的增删改、交差合、拆分等。但List不是专业的结构化数据对象,一旦涉及字段结构相关的功能,Kotlin就很难实现了。比如,取Orders中的两个字段组成新的结构化数据对象。 data class CliAmt(var Client: String, var Amount: Double) var CliAmts=Orders.map{it.let{CliAmt(it.Client,it.Amount) }} 上面的功能很常用,相当于简单SQL语句select Client,Amount from Orders,但Kotlin写起来就很繁琐,不仅要事先定义新结构,还要硬编码完成字段的赋值。简单的取字段功能都这么繁琐,高级些的功能就更麻烦了,比如:按字段序号取、按参数取、获得字段名列表、修改字段结构、在字段上定义键和索引、按字段查询计算。 Scala也有List,与Kotlin区别不大,但Scala为结构化数据处理设计了更加专业的数据对象DataFrame(以及RDD、DataSet)。 DataFrame是有结构的数据流,与数据库结果集有些相似,都是无序集合,因此不支持按下标取数,只能变相实现。比如,第10条记录: Orders.limit(10).tail(1)(0) 可以想象,凡与顺序相关的计算,DataFrame实现起来都比较麻烦,比如区间、移动平均、倒排序等。 除了数据无序,DataFrame也不支持修改(immutable特性),如果想改变数据或结构,必须生成新的DataFrame。比如修改字段名,实际上要通过复制记录来实现: Orders.selectExpr("Client as Cli") DataFrame支持常见的集合计算,比如拆分、合并、交差合并,其中并集可通过合集去重实现,但因为要通过复制记录来实现,集合计算的性能普遍不高。 虽然有不少缺点,但DataFrame是专业的结构化数据对象,字段访问方面的能力是Kotlin无法企及的。比如,获得元数据/字段名列表: Orders.schema.fields.map(it=>it.name).toList 还可以方便地用字段取数,比如,取两个字段形成新dataframe: Orders.select("Client","Amount") //可以只用字段名 或用计算列形成新DataFrame: Orders.select(Orders("Client"),Orders("Amount")+1000) //不能只用字段名 遗憾的是,DataFrame只支持用字符串形式的名字来引用字段,不支持用字段序号或默认名字,导致很多场景下不够方便。此外,DataFrame也不支持定义索引,无法进行高性能随机查询,专业性还有缺陷。 SPL的结构化数据对象是序表,优点是足够专业,简单易用,表达能力强。 按序号访问成员: Orders(3) //按下标取记录,从1开始 Orders.to(3) //前3条记录 Orders.m(1,3,5,7:10) //序号是1、3、5、7-10的记录 按倒数序号取记录,独特之处在于支持负号表示倒数,比Kotlin专业且方便: Orders.m(-1,-3,-5) //倒数第1,3,5条 Orders.m(1,-1) //第1条和最后1条 1 2 作为集合的一种,序表也支持集合成员的增删改、交并差合、拆分等功能。由于序表和List一样都是可变集合(mutable),集合计算时尽可能使用游离记录,而不是复制记录,性能比Scala好得多,内存占用也少。 序表是专业的结构化数据对象,除了集合相关功能外,更重要的是可以方便地访问字段。比如,获得字段名列表: Orders.fname() 1 取两个字段形成新序表: Orders.new(Client,Amount) 1 用计算列形成新序表: Orders.new(Client,Amount*0.2) 1 修改字段名: Orders.alter(;OrderDate) //不复制记录 1 有些场景需要用字段序号或默认名字访问字段,SPL都提供了相应的访问方法: Orders(Client) //按字段名(表达式取) Orders([#2,#3]) //按默认字段名取 Orders.field(“Client”) //按字符串(外部参数) Orders.field(2) //按字段序号取 1 2 3 4 作为专业的结构化数据对象,序表还支持在字段上定义键和索引: Orders.keys@i(OrderID) //定义键,同时建立哈希索引 Orders.find(47) //用索引高速查找 1 2 计算函数 Kotlin支持部分基本计算函数,包括:过滤、排序、去重、集合的交叉合并、各类聚合、分组汇总。但这些函数都是针对普通集合的,如果计算目标改成结构化数据对象,计算函数库就显得非常不足,通常就要辅以硬编码才能实现计算。还有很多基本的集合运算是Kotlin不支持的,只能自行编码实现,包括:关联、窗口函数、排名、行转列、归并、二分查找等。其中,归并和二分查找等属于次序相关的运算,由于Kotlin List是有序集合,自行编码实现这类运算不算太难。总体来讲,面对结构化数据计算,Kotlin的函数库可以说较弱。 Scala的计算函数比较丰富,且都是针对结构化数据对象设计的,包括Kotlin不支持的函数:排名、关联、窗口函数、行转列,但基本上还没有超出SQL的框架。也有一些基本的集合运算是Scala不支持的,尤其是与次序相关的,比如归并、二分查找,由于Scala DataFrame沿用了SQL中数据无序的概念,即使自行编码实现此类运算,难度也是非常大的。总的来说,Scala的函数库比Kotlin丰富,但基本运算仍有缺失。 SPL的计算函数最丰富,且都是针对结构化数据对象设计的,SPL极大地丰富了结构化数据运算内容,设计了很多超出SQL的内容,当然也是Scala/Kotlin不支持的函数,比如有序计算:归并、二分查找、按区间取记录、符合条件的记录序号;除了常规等值分组,还支持枚举分组、对齐分组、有序分组;将关联类型分成外键和主子;支持主键以约束数据,支持索引以快速查询;对多层结构的数据(多表关联或Json\XML)进行递归查询等。 以分组为例,除了常规的等值分组外,SPL还提供了更多的分组方案: 枚举分组:分组依据是若干条件表达式,符合相同条件的记录分为一组。 对齐分组:分组依据是外部集合,记录的字段值与该集合的成员相等的分为一组,组的顺序与该集合成员的顺序保持一致,允许有空组,可单独分出一组“不属于该集合的记录”。 有序分组:分组依据是已经有序的字段,比如字段发生变化或者某个条件成立时分出一个新组,SPL直接提供了这类有序分组,在常规分组函数上加个选项就可以完成,非常简单而且运算性能也更好。其他语言(包括SQL)都没有这种分组,只能费劲地转换为传统的等值分组或者自己硬编码实现。 下面我们通过几个常规例子来感受一下这三种语言在计算函数方式的差异。 排序 按Client顺序,Amount逆序排序。Kotlin: Orders.sortedBy{it.Amount}.sortedByDescending{it.Client} 1 Kotlin代码不长,但仍有不便之处,包括:逆序正序是两个不同的函数,字段名必须带表名,代码写出的字段顺序与实际的排序顺序相反。 Scala: Orders.orderBy(Orders("Client"),-Orders("Amount")) 1 Scala简单多了,负号代表逆序,代码写出的字段顺序与排序的顺序相同。遗憾之处在于:字段仍要带表名;编译型语言只能用字符串实现表达式的动态解析,导致代码风格不统一。 SPL: Orders.sort(Client,-Amount) 1 SPL代码更简单,字段不必带表名,解释型语言代码风格容易统一。 分组汇总 Kotlin: data class Grp(var Dept:String,var Gender:String) data class Agg(var sumAmount: Double,var rowCount:Int) var result1=data.groupingBy{Grp(it!!.Dept,it.Gender)} .fold(Agg(0.0,0),{acc, elem -> Agg(acc.sumAmount + elem!!.Amount,acc.rowCount+1)}) .toSortedMap(compareBy { it.Dept }.thenBy { it.Gender }) Kotlin代码比较繁琐,不仅要用groupingBy和fold函数,还要辅以硬编码才能实现分组汇总。当出现新的数据结构时,必须事先定义才能用,比如分组的双字段结构、汇总的双字段结构,这样不仅灵活性差,而且影响解题流畅性。最后的排序是为了和其他语言的结果顺序保持一致,不是必须的。 Scala: val result=data.groupBy(data("Dept"),data("Gender")).agg(sum("Amount"),count("*")) 1 Scala代码简单多了,不仅易于理解,而且不用事先定义数据结构。 SPL: data.groups(Dept,Gender;sum(Amount),count(1)) 1 SPL代码最简单,表达能力不低于SQL。 关联计算 两个表有同名字段,对其关联并分组汇总。Kotlin代码: data class OrderNew(var OrderID:Int ,var Client:String, var SellerId:Employee ,var Amount:Double ,var OrderDate:Date ) val result = Orders.map { o->var emp=Employees.firstOrNull{it.EId==o.EId} emp?.let{ OrderNew(o.OrderID,o.Client,emp,o.Amount,o.OrderDate)} } .filter {o->o!=null} data class Grp(var Dept:String,var Gender:String) data class Agg(var sumAmount: Double,var rowCount:Int) var result1=data.groupingBy{Grp(it!!.EId.Dept,it.EId.Gender)} .fold(Agg(0.0,0),{acc, elem -> Agg(acc.sumAmount + elem!!.Amount,acc.rowCount+1)}) .toSortedMap(compareBy { it.Dept }.thenBy { it.Gender }) Kotlin代码很繁琐,很多地方都要定义新数据结构,包括关联结果、分组的双字段结构、汇总的双字段结构。 Scala val join=Orders.as("o").join(Employees.as("e"),Orders("EId")===Employees("EId"),"Inner") val result= join.groupBy(join("e.Dept"), join("e.Gender")).agg(sum("o.Amount"),count("*")) Scala比Kolin简单多了,不用繁琐地定义数据结构,也不必硬编码。 SPL更简单: join(Orders:o,SellerId;Employees:e,EId).groups(e.Dept,e.Gender;sum(o.Amount),count(1)) 1 综合数据处理对比 CSV内容不规范,每三行对应一条记录,其中第二行含三个字段(即集合的集合),将该文件整理成规范的结构化数据对象,并按第3和第4个字段排序. Kotlin: data class Order(var OrderID: Int,var Client: String,var SellerId: Int, var Amount: Double, var OrderDate: Date) var Orders=ArrayList() var sdf = SimpleDateFormat("yyyy-MM-dd") var raw=File("d:\\threelines.txt").readLines() raw.forEachIndexed{index,it-> if(index % 3==0) { var f234=raw[index+1].split("\t") var r=Order(raw[index].toInt(),f234[0],f234[1].toInt(),f234[2].toDouble(), sdf.parse(raw[index+2])) Orders.add(r) } } var result=Orders.sortedByDescending{it.Amount}.sortedBy{it.SellerId} Koltin在数据处理方面专业性不足,大部分功能要硬写代码,包括按位置取字段、从集合的集合取字段。 Scala: val raw=spark.read.text("D:/threelines.txt") val rawrn=raw.withColumn("rn", monotonically_increasing_id()) var f1=rawrn.filter("rn % 3==0").withColumnRenamed("value","OrderId") var f5=rawrn.filter("rn % 3==2").withColumnRenamed("value","OrderDate") var f234=rawrn.filter("rn % 3==1") .withColumn("splited",split(col("value"),"\t")) .select(col("splited").getItem(0).as("Client") ,col("splited").getItem(1).as("SellerId") ,col("splited").getItem(2).as("Amount")) f1.withColumn("rn1",monotonically_increasing_id()) f5=f5.withColumn("rn1",monotonically_increasing_id()) f234=f234.withColumn("rn1",monotonically_increasing_id()) var f=f1.join(f234,f1("rn1")===f234("rn1")) .join(f5,f1("rn1")===f5("rn1")) .select("OrderId","Client","SellerId","Amount","OrderDate") val result=f.orderBy(col("SellerId"),-col("Amount")) Scala在数据处理方面更加专业,大量使用结构化计算函数,而不是硬写循环代码。但Scala缺乏有序计算能力,相关的功能通常要添加序号列再处理,导致整体代码冗长。 SPL: A 1 =file("D:\\data.csv").import@si() 2 =A1.group((#-1)\3) 3 =A2.new(~(1):OrderID, (line=~(2).array("\t"))(1):Client,line(2):SellerId,line(3):Amount,~(3):OrderDate ) 4 =A3.sort(SellerId,-Amount) SPL在数据处理方面最专业,只用结构化计算函数就可以实现目标。SPL支持有序计算,可以直接按位置分组,按位置取字段,从集合中的集合取字段,虽然实现思路和Scala类似,但代码简短得多。 应用结构 Java应用集成 Kotlin编译后是字节码,和普通的class文件一样,可以方便地被Java调用。比如KotlinFile.kt里的静态方法fun multiLines(): List,会被Java正确识别,直接调用即可: java.util.List result=KotlinFileKt.multiLines(); result.forEach(e->{System.out.println(e);}); Scala编译后也是字节码,同样可以方便地被Java调用。比如ScalaObject对象的静态方法def multiLines():DataFrame,会被Java识别为Dataset类型,稍做修改即可调用: org.apache.spark.sql.Dataset df=ScalaObject.multiLines(); df.show(); SPL提供了通用的JDBC接口,简单的SPL代码可以像SQL一样,直接嵌入Java: Class.forName("com.esproc.jdbc.InternalDriver"); Connection connection =DriverManager.getConnection("jdbc:esproc:local://"); Statement statement = connection.createStatement(); String str="=T(\"D:/Orders.xls\").select(Amount>1000 && Amount<=3000 && like(Client,\"*s*\"))"; ResultSet result = statement.executeQuery(str); 的SPL代码可以先存为脚本文件,再以存储过程的形式被Java调用,可有效降低计算代码和前端应用的耦合性。 Class.forName("com.esproc.jdbc.InternalDriver"); Connection conn =DriverManager.getConnection("jdbc:esproc:local://"); CallableStatement statement = conn.prepareCall("{call scriptFileName(?, ?)}"); statement.setObject(1, "2020-01-01"); statement.setObject(2, "2020-01-31"); statement.execute(); SPL是解释型语言,修改后不用编译即可直接执行,支持代码热切换,可降低维护工作量,提高系统稳定性。Kotlin和Scala是编译型语言,编译后必须择时重启应用。 开源是每个程序员的精神,大家可以一起交流学习 交互式命令行 Kotlin的交互式命令行需要额外下载,使用Kotlinc命令启动。Kotlin命令行理论上可以进行任意复杂的数据处理,但因为代码普遍较长,难以在命令行修改,还是更适合简单的数字计算: >>>Math.sqrt(5.0) 2.236.6797749979 Scala的交互式命令行是内置的,使用同名命令启动。Scala命令行理论上可以进行数据处理,但因为代码比较长,更适合简单的数字计算: scala>100*3 rest1: Int=300 SPL内置了交互式命令行,使用“esprocx -r -c”命令启动。SPL代码普遍较短,可在命令行进行简单的数据处理。 (1): T("d:/Orders.txt").groups(SellerId;sum(Amount):amt).select(amt>2000) (2):^C D:\raqsoft64\esProc\bin>Log level:INFO 1 4263.900000000001 3 7624.599999999999 4 14128.599999999999 5 26942.4 ———————————————— 版权声明:本文为CSDN博主「程序员springmeng」的原创文章,遵循CC 4.0 BY-SA版权协议,转载请附上原文出处链接及本声明。 原文链接:https://blog.csdn.net/mengchuan6666/article/details/126616736
-
近日,在 “2022 IDC中国数字金融论坛”上,国际权威咨询机构IDC联合蚂蚁集团正式发布了《十大风控技术趋势指南》白皮书。这是风控行业技术创新的一次风向标,也意味着和黑灰产对抗中技术升级迫在眉睫。当今的商业模式已不同于往昔,随着数字化进程的进一步加快,金融机构必须要时刻为可能出现的业务风险做好准备。面对正在走向无边界和强对抗的新型重大风险,金融机构如何与之博弈,并始终领先一步?这正是“ IDC《十大风控技术趋势指南》”将深入探讨的议题。数字支付激增 新型风险类型相伴相生新冠疫情算得上数字化发展的一个“加速器”,但其实早在疫情出现之前,数字服务领域已经有了大规模的转型:从线上互动,到数字支付,再到依托于数字平台而生的新服务。疫情的出现加速了这一趋势,加速的势头预计将保持到2030年。图1数据显示,2020年到2025年,全球消费者数字支付市场预计增长2.2倍,而在2025到2030年期间,上涨幅度预计将进一步增至3.4倍。数字化世界的机遇和潜力巨大,但也充满了风险。随着企业加速运营调整以应对数字化进程,这种激增的趋势带来了明显的合规风险和业务风险,让黑产有机可乘。IDC的一项研究表明,相较于2020年,在2021年,亚太地区52%的企业因遭遇诈骗而蒙受的损失上涨了至少5%,26%的企业损失上涨了至少11%。由于支付管控不力及降低风险手段的不到位,欺诈活动现在越来越猖獗。黑灰产的欺诈手法正在不断升级,欺诈套路也变得越来越复杂。值得关注的十大新技术能力面对快速变化的欺诈发展形势,传统降低风险的做法和欺诈检测工具是否能及时应对?如果不采用新的工具和技术,企业是否能够安然地扩大其数字化业务的规模?现有的基础设施是否能够支撑企业分析海量数据、检测欺诈,尤其是新型欺诈?对于亚太地区的银行、商家、支付公司和其他金融机构来说,这些问题的答案可能都是否定的。本节中,我们将重点说明十项科技趋势,凭借这些能力,金融机构才能够有机会实现可信的智能黑灰产对抗。01 人工智能,风控能力提升的基础预计到2025年,银行业还将再投入约310亿美元用于在现有系统中嵌入人工智能技术。在接受调查的100位来自全球银行业的高管中,多数人表示他们会将欺诈管理作为重点,其中,有些银行在与欺诈相关的场景用例中已经应用了人工智能,包括开户欺诈(57%)、支付欺诈检测(57%)、欺诈操作和调查(53%)还有反洗钱监测(46%)。自2022年开始,人工智能将成为打击欺诈活动的一个重要基础能力,人工智能将有效缩短决策时间,在7*24小时的全天候业务中,帮助实现客户快捷、无缝的交易体验,同时确保决策的准确性。02 威胁情报的挖掘技术, 为风险防控提供有效依据IT安全解决方案不胜枚举,而市场仍然对多种威胁情报有强烈需求。因此基于巡检技术的胁情报挖掘和分享能够持续为金融机构提供与新型威胁、欺诈迹象相关的信息。金融机构需要对这些情报进行审查,同时记录不同威胁情报效率和准确度得分,以便更好地了解不同的线索,进而指导对各种威胁的检测、识别、调查和处理。黑灰产通常不会只在一个平台犯案,因此威胁情报对金融机构来说至关重要,例如亚太地区的许多银行协会,他们会定期分享他们感知到的威胁情报,并和行业分享应对举措。对金融机构来说,你得到的情报越多越准确,就越有可能在风险防控中领先于黑灰产。03 全图风控, 实现动态可视事实风险挖掘金融风险决策是一个不断对抗升级的过程,从单一事件和孤立行为来分析无法获得准确决策。随着大规模图计算技术的发展,风险防控将从单一时间切片的图数据,走向基于时序的图数据,该防控方式将有效沉淀如账户盗用、电信网络诈骗、套利等风险特征,通过知识表征推理发现更多稀薄关系和隐藏风险,结合规则推理、规则挖掘与规则学习挖掘更多风险模式并有效泛化,让风险知识和实时交易事件联动实现动态图推理,形成全局的洞察,构建实时监控体系。基于大规模图技术的全图风控能够支持千亿级的金融风险知识图谱,进而为管理者们提供全面、可见、动态、实时的交易风险概览,使他们能够监测风险并及时决策。04 高效的算力体系,为精准流畅风险防控提供算力支撑交易和互动的数量、频率都在急剧增长,随之而来的是数据的激增,企业几乎要被海量的数据所淹没。此外,消费者设备、支付渠道、5G网络、物联网等也在不断产生新的数据。现在,企业所面临的挑战是通过分析从不同来源(结构化和非结构化)收集到的数据,从而发现欺诈的线索。然而,生成的大量数据可能会使存储和处理的环节负担过重,进而让不法分子有机可乘,组织比如跨境洗钱、非法交易等网络犯罪。一旦处理和分析数据的机制存在缺陷的话,那么虚假交易的中间人就很可能“隐身”其中,为了能够实时、准确地检测到欺诈行为,只有将传统的架构转为云计算和多节点高效算力体系,才能利用更高的计算效率来支撑人工智能/机器学习的计算需求。05 极速风控,实现更快的实时风险决策欺诈检测的实效性对金融机构来说至关重要,分析决策环节的每一秒延时都会降低用户体验,也让金融机构和用户增加一份资损的风险。在登录、交易支付、验证检查或用户验证等环节,实时决策的能力有赖于风险情报的收集和风控系统强大的分析和计算能力,而如何解决大规模风险数据计算中的耗时问题是行业面临的一大挑战。极速风控通过预测的方式将风险识别和风险决策进行解耦,通过提前风险计算,提高决策时的风险判断效率,实现毫秒级的实时风险决策。06 主动式风控, 在即时响应基础上主动出击传统的风险管理解决方案大都是被动的“事后应对”:即在不利事件发生后,基于已有信息做出判断,采取保护性行动,以便之后能够及时应对类似的攻击。但这还远远不够,尤其是面对技术越来越好、作案手段不断演进的欺诈团伙。随着人工智能及其相关技术的发展,企业主动应对潜在的风险变得可能,例如通过主动和用户产生交互,来获得更多的风险信息,帮助平台做更好的风险判断,同时给到用户更好的安全服务。以本人授权的被诈骗支付为例,传统的风险管理系统仅能在检测到风险后限制或冻结交易;而现在,系统能够在发现潜在风险后,以图文提示、电话等多模态交互方式进一步确认风险,提醒用户主动意识到欺诈风险。07 端云协同,提高计算效能保护用户隐私随着企业越来越重视隐私保护和用户体验,传统的风险防控将面临全新的挑战,为了应对隐私保护和用户体验的挑战,端云协同的方案应运而生。受海量流媒体数据的驱动,企业需要让数据处理环节更靠近数据的来源,以进一步降低延迟、加快决策,减少个人数据的传输。通过端云协同的风控方案,企业可以让隐私数据计算在用户智能终端(如手机)中进行,将不含隐私信息的决策结果输送到云端,以实现“端云协同”的风控保障。08 多方风控,确保安全的跨机构协作数字化世界愈加互联互通,但很多时候,即使一家公司内的风险数据都没有被整合,更不用说行业间风险数据的互联互通。基于此,多方风控技术已在广泛试点使用,不但让多方在共同应对欺诈时实现数据、模型和分析结果的共享,而无需牺牲数据隐私或数据集的质量。有了这一更高效的协作方式,多方均可提升自身在鉴别和应对风险方面的能力。多方风控主要由区块链及隐私计算技术支撑,比如可信执行环境(TEE), 多方安全计算和联邦学习,使得不同的机构能够在数据隐私得到极好保护的前提下进行风险数据共享,甚至联合建模。因此,为应对连通性风险,各商家、银行和第三方支付机构之间的“互联互通”十分必要,同时,还须保证这种“互联互通” 的安全性。09 可信AI,智能风控系统的安全基础人工智能(AI)的应用是风险管理中出现的新常态。但是,AI不仅仅可以为好人所用,也可以被黑灰产作为突破口,或者攻击武器。由于风险防控是一场和犯罪团伙的竞速赛,企业必须开始考虑他们以人工智能驱动的风险防控系统是否足够稳健、可靠,能够扛得住黑灰产的攻击。这时对抗智能就变得尤为重要,它建立在经济学的博弈论框架之上,通过模拟攻击者和防御者之间的冲突,让机器自动且实时、动态地对自身系统进行安全性攻击,从而提升模型能力,使模型更加鲁棒(robust),处理结果更加准确。先进的欺诈管理解决方案已经采用了对抗智能技术,以提升人工智能模型的稳健性,此涉及的技术很多,包括像防御性的对抗性权重扰动(AWP)、投影梯度下降(PGD)等概念和技术。在智能数字化服务中,我们必须尽可能地严格看待人工智能/ 机器学习模型所做出的决策。如果人工智能/ 机器学习的决策是基于不完整、低质量、非客观的数据集,通过错误的建模方式和错误的变量集而做出的,在未来可能会引发了诸多争议。因此,企业应当建立一个值得信赖、可靠且可追溯的AI安全框架,来更好地管理AI相关风险。尽管AI构成了应对复杂欺诈案件的解决方案,但如果没有恰当的AI治理框架,AI也可能会影响用户体验甚至是破坏品牌声誉。AI模型的安全性也需要保护,因为它们也可能会被那些有技术团队的专业作案团伙所破坏。10 用户行为分析(UBA),将变得愈加重要金融机构在行为分析方面的投资正在逐年增加,以提升其分析客户资料、互动模式和交易数据的能力,此外,行为分析还能帮助银行发现可疑活动,检测和预防欺诈。多年以来,在IT安全市场上,人们都是在不利事件发生后才想起这一能力,因而直到现在,UBA的相关投资仍相对缺乏;但是,未来对UBA的投资估计不会小。当然,分析的本质决定了对其投入的时间越多,效率越会提升。要想实现有效的UBA,需要花费大量的时间并进行多次的细微调整,同时还需制定一条恰当的路线图。随着时间的推移,企业使用UBA会愈加成熟,逐渐形成自己的反馈回路并获得一系列的结果,根据这些结果,他们可以再进行建模。来源:IDC、蚂蚁集团、蚂蚁技术AntTech
-
北京,2022年8月30日IDC近日发布了《2022年V2全球大数据支出指南》(IDC Worldwide Big Data and Analytics Spending Guide)。本次发布新增2026年预测数据,从技术、行业、企业规模等维度发掘未来五年(2021-2026)全球大数据市场中的发展趋势和潜在机会,同时对2021年的市场情况进行了梳理。大数据市场概览IDC数据显示,2021年全球大数据市场的IT总投资规模为2,176.1亿美元,并有望在2026年增至4,491.1亿美元,五年预测期内(2021-2026)实现约15.6%的复合增长率(CAGR)。聚焦中国市场,IDC预计,2026年中国大数据IT支出规模预计为359.5亿美元,市场规模位列单体国家第二。从增速的角度来看,中国大数据IT支出五年CAGR约为21.4%,位列全球第一。中国大数据市场增速持续领跑全球,呈现出强劲的增长态势,市场前景广阔。随着数字经济、数字化转型、新基建等投资建设进一步加快,中国终端用户对大数据硬件、软件、服务的需求将稳步扩大。 技术维度IDC预测,到2026年,中国大数据硬件市场IT投资规模将达到137.2亿美元,超过2021年投资规模的两倍。值得关注的是,未来五年,硬件市场仍将是中国大数据市场占比最高的一级子市场,占比规模接近四成。聚焦中国大数据软件市场,2026年大数据软件将成为第二大技术市场。大数据软件以26.9%的五年CAGR强势增长,软件IT投资规模逐年接近硬件市场。其中,人工智能软件平台(AI Software Platforms)市场和终端用户查询、报告和分析(End-User Query, Reporting and Analysis Tools)市场将主导中国大数据软件IT投资,两者共计近软件投资总规模的四成。从增速的角度来看,内容分析(Content Analytics Tools)技术子市场增速亮眼,该市场将以41.1%的五年CAGR快速扩大规模。未来大数据软件市场将发挥承上启下的关键作用,与上下游产品形成耦合榫卯结构。从中国大数据服务市场的角度来看,2026年中国大数据服务市场规模将接近百亿大关。面对全球服务市场增速放缓的大趋势,中国大数据服务市场将以略高于全球平均水平的五年CAGR稳步增长。行业应用从行业终端用户的角度来看,至2026年,专业服务、电信、金融和政府将成为大数据相关IT支出的主力行业。具体而言,专业服务、电信、银行和地方政府将会贡献超过50%的中国大数据IT投资。就增速而言,医疗保健行业将以30.9%的五年CAGR成为增长最快的行业终端用户。此外,专业服务、离散制造、电信等行业也展现出了较大的发展潜力。IDC调研显示,各行业领域企业都在不断探索布局大数据处理分析产品和完整解决方案,文娱、电商、社交等多样式互联网产品服务创新将持续带动市场增长,对于信息化基础较好、数据就绪度较高、市场服务要求更高的电信、金融等行业,平台管理、决策分析等组合型产品市场前景明朗。另外,在政府专项政策推动下,智慧城市、智能制造、智慧医疗、智慧农业等领域也将迎来新的机遇。终端用户企业规模IDC《全球大数据支出指南》将终端用户企业规模由上至下分为了五个区间,从企业规模的维度对大数据支出情况做出进一步透视。和上一版相比,中国大数据市场集中度有所提高。雇员超过1,000人的超大型企业在五年预测期内(2021-2026)占据整个中国市场支出的65%左右,较之前预测小幅上调。中小型企业整体增速较快,但市场占比较小。大数据市场呈现横纵一体化发展格局,大型及超大型企业依托优势数据资源、丰富行业场景经验、高水平信息化技术、规模化服务体系以及较强预算支出能力,表现出较强竞争力。而中小型企业依托核心技术和业务资源优势,在专业化需求上也拥有一定优势。
-
10:12 Cannot download sources Sources not found for: org.apache.spark:spark-hive_2.11:2.4.5-hw-ei-302002
-
摘要:SPL实现了更优算法,性能远远超过存储过程,能显著提高单机计算效率,非常适合跑批计算。本文分享自华为云社区《Java开源专业计算引擎:跑批真的这么难吗?》,作者: Java李杨勇。业务系统产生的明细数据通常要经过加工处理,按照一定逻辑计算成需要的结果,用以支持企业的经营活动。这类数据加工任务一般会有很多个,需要批量完成计算,在银行和保险行业常常被称为跑批,其它像石油、电力等行业也经常会有跑批的需求。大部分业务统计都会要求以某日作为截止点,而且为了不影响生产系统的运行,跑批任务一般会在夜间进行,这时候才能将生产系统当天产生的新明细数据导出来,送到专门的数据库或数据仓库完成跑批计算。第二天早上,跑批结果就可以提供给业务人员使用了。和在线查询不同,跑批计算是定时自动执行的离线任务,不会出现多人同时访问一个任务的情况,所以没有并发问题,也不必实时返回结果。但是,跑批必须在规定的窗口时间内完成。比如某银行的跑批窗口时间是晚上8:00到第二天早上7:00,如果到了早上7:00跑批任务还没有完成,就会造成业务人员无法正常工作的严重后果。跑批任务涉及的数据量非常大,很可能用到所有的历史数据,而且计算逻辑复杂、步骤众多,所以跑批时间经常是以小时计的,一个任务两三小时是家常便饭,跑到十个小时也不足为奇。随着业务的发展,数据量还在不断增加。跑批数据库的负担快速增长,就会发生整晚都跑不完的情况,严重影响用户的业务,这是无法接受的。问题分析要解决跑批时间过长的问题,必须仔细分析现有的系统架构中的问题。跑批系统比较典型的架构大致如下图:从图上看,数据要从生产数据库取出,存入跑批数据库。跑批数据库通常是关系型的,编写存储过程代码完成跑批计算。跑批的结果一般不会直接使用,而是再从跑批数据库中导出,采用接口文件的方式提供给其他系统,或者再导入其他系统数据库。这是比较典型的架构,图中的生产数据库也可能是某个中央数据仓库或者Hadoop等。一般情况下,生产库和跑批库不会是同一种数据库,它们之间往往通过文件的方式传递数据,这样也比较有利于降低耦合度。跑批计算完成后,结果要给多个应用系统使用,一般也都是以文件方式传递。跑批很慢的第一个原因,是用来完成跑批任务的关系数据库入库、出库太慢。由于关系数据库的存储和计算能力具有封闭性,数据的进出要做过多的约束检查和安全处理,当数据量较大时,写入读出的效率非常低,耗时会非常长。所以,跑批数据库导入文件数据的过程,以及跑批计算结果再导出文件的过程都会很慢。跑批很慢的第二个原因,是存储过程性能差。由于SQL的语法体系过于陈旧,存在诸多限制,很多高效的算法无法实施,所以存储过程中的SQL语句计算性能很不理想。而且,业务逻辑比较复杂的时候很难用一个SQL实现,经常要分成多个步骤,用十几甚至几十个SQL语句才能完成。每个SQL的中间结果,都要存入临时表给后续步骤的SQL使用。临时表数据量较大时就必须落地,会造成大量的数据写出。而数据库的写出要比读入性能差很多,会严重拖慢整个存储过程。对于更复杂的计算,甚至很难用SQL语句直接实现,需要用数据库游标遍历取出数据,循环计算。但数据库游标遍历计算性能又要比SQL语句差很多,一般也都不直接支持多线程并行计算,很难利用多CPU核的计算能力,会让计算性能更加糟糕。那么,是否可以考虑用分布式数据库来代替传统关系数据库,通过增加节点数量的办法,来提高跑批任务的速度呢?答案仍然是不可行。主要原因是跑批计算的逻辑相当复杂,即使是用传统数据库的存储过程,也常常要写几千甚至上万行代码,而分布式数据库的存储过程计算能力还比较弱,很难实现这么复杂的跑批计算。而且,当复杂计算任务不得不分成多个步骤时,分布式数据库也面临中间结果落地的问题。由于数据可能在不同的节点上,所以前序步骤将中间结果落地,后续步骤再读取的时候,都会造成大量跨网络的读写操作,性能很不可控。这时,也不能采用分布式数据库依靠数据冗余来提升查询速度的办法。这是因为,查询之前可以预先准备好多份冗余数据,但是,跑批的中间结果是临时生成的,如果冗余的话就要临时生成多份,整体的性能只会变得更慢。所以,现实的跑批业务通常仍然是使用大型单体数据库进行,计算强度太大时会采用类似ExaData这样的一体机(ExaData是多数据库,但被Oracle专门优化过,可以看成是个超大型单体数据库)。虽然很慢,但是暂时找不到更好的选择,只有这类大型数据库有足够的计算能力,所以只能用它来完成跑批任务了。SPL用于跑批开源的专业计算引擎SPL提供了不依赖数据库的计算能力,直接利用文件系统计算,可以解决关系数据库出库入库太慢的问题。而且SPL实现了更优算法,性能远远超过存储过程,能显著提高单机计算效率,非常适合跑批计算。利用SPL实现的跑批系统新架构是下面这样的:在新架构中,SPL解决了造成跑批慢的两大瓶颈问题。首先来看数据的入库、出库问题。SPL可以直接基于生产库导出的文件计算,不必再将数据导入到关系数据库中。完成跑批计算后,SPL还能将最终结果直接存储成文本文件等通用格式,传递给其他应用系统,避免了原有跑批数据库的出库操作。这样一来,SPL就省去了关系数据库缓慢的入库、出库过程。下面再来看计算的过程。SPL提供了更优的算法(有许多是业界首创),计算性能远远超过存储过程和SQL语句。这些高性能算法包括:这些高性能算法可以应用于跑批任务中的常见JOIN计算、遍历、分组汇总等,能有效提升计算速度。例如,跑批任务常常要遍历整个历史表。有些情况下,对一个历史表还要遍历好多次,来完成多种业务逻辑的计算。历史表数据量一般都很大,每次遍历都要消耗很多的时间。此时我们可以应用SPL的遍历复用机制,仅对大表遍历一次,就可以同时完成多种计算,可以节省大量时间。SPL的多路游标能做到数据的并行读取和计算,即使是很复杂的跑批逻辑,也可以利用多CPU核实现多线程并行运算。而数据库游标是很难并行的,这样一来,SPL的计算速度常常可以达到存储过程的数倍。SPL的延迟游标机制,可以在一个游标上定义多个计算步骤,之后让数据流按顺序依次完成这些步骤,实现链式计算,能够有效减少中间结果落地的次数。在数据必须落地的情况下,SPL也可以将中间结果存成内置的高性能数据格式,供下一个步骤使用。SPL高性能存储基于文件,采用有序压缩存储、自由列式存储、倍增分段、自有压缩编码等技术,减少了硬盘占用,读写速度要远远好于数据库。应用效果SPL在技术架构上打破了关系型跑批数据库存在的两大瓶颈,在实际应用中也取得了非常好的效果。L 银行跑批任务采用传统架构,以关系数据库作为跑批数据库,用存储过程编程实现跑批逻辑。其中,贷款协议存储过程需要执行 2 个小时,而且是很多其他跑批任务的前序任务,耗时这么久,对整个跑批任务造成了严重影响。采用SPL后,使用高性能列存、文件游标、多线程并行、小结果内存分组、游标复用等高性能算法和存储机制,将原来2个小时的计算时间缩短为10分钟,性能提高12倍。而且,SPL代码更简洁。原存储过程3300多行,改为SPL后,仅有500格语句,代码量减少了6倍多,大大提高了开发效率。P保险公司的车险业务中,需要用往年历史保单来关联新的保单,在跑批中称为历史保单关联任务。原来也采用关系数据库完成跑批,存储过程计算10天的新增保单关联历史保单,运行时间47分钟;30天则需要112分钟,接近2小时;如果日期跨度更大,运行时间就会长的无法忍受,基本就变成不可能完成的任务了。采用SPL后,应用了高性能文件存储、文件游标、有序归并分段取出、内存关联和遍历复用等技术,计算10天新增保单仅需13分钟;30天新增保单只需要17分钟,速度提高了近7倍。而且,新算法执行的时间随着保单天数的增长并不是很大,并没有像存储过程那样成正比的增长。从代码总量来看,原来存储过程有2000行代码,去掉注释后还有1800多行,而SPL的全部代码只有不到500格,不到原来的1/3。T银行通过互联网渠道发放贷款的明细数据,需要每天执行跑批任务,统计汇总指定日期之前的所有历史数据。跑批任务采用关系数据库的SQL语句实现,运行总时间7.8小时,占用了过多的跑批时间,甚至影响了其他的跑批任务,必须优化。采用SPL后,应用了高性能文件、文件游标、有序分组、有序关联、延迟游标、二分法等技术,原来需要7.8小时的跑批任务,单线程仅需180秒,2线程仅需137秒,速度提高了204倍。
-
鲲鹏920服务器,1台NameNode, 3台DataNode,Hadoop 3.3.1 使用官网 aarch64包 cid:link_0works配置正确,副本数设置为3.进行dfs write测试时会报如下错误(使用Hadoop 2.6不会),求解决方法: 2022-08-25 10:21:38,080 INFO org.apache.hadoop.hdfs.StateChange: BLOCK* allocate blk_1073743798_2974, replicas=10.10.10.4:9866, 10.10.10.3:9866 for /benchmarks/TestDFSIO/io_control/in_file_test_io_1973 2022-08-25 10:21:38,082 INFO org.apache.hadoop.hdfs.StateChange: DIR* completeFile: /benchmarks/TestDFSIO/io_control/in_file_test_io_1973 is closed by DFSClient_NONMAPREDUCE_-1180995040_1 2022-08-25 10:21:38,083 INFO org.apache.hadoop.hdfs.server.blockmanagement.BlockPlacementPolicy: Not enough replicas was chosen. Reason: {NO_REQUIRED_STORAGE_TYPE=1} 2022-08-25 10:21:38,083 INFO org.apache.hadoop.hdfs.server.blockmanagement.BlockPlacementPolicy: Not enough replicas was chosen. Reason: {NO_REQUIRED_STORAGE_TYPE=1} 2022-08-25 10:21:38,083 INFO org.apache.hadoop.hdfs.server.blockmanagement.BlockPlacementPolicy: Not enough replicas was chosen. Reason: {NO_REQUIRED_STORAGE_TYPE=1} 2022-08-25 10:21:38,083 WARN org.apache.hadoop.hdfs.server.blockmanagement.BlockPlacementPolicy: Failed to place enough replicas, still in need of 1 to reach 3 (unavailableStorages=[], storagePolicy=BlockStoragePolicy{HOT:7, storageTypes=[DISK], creationFallbacks=[], replicationFallbacks=[ARCHIVE]}, newBlock=true) For more information, please enable DEBUG log level on org.apache.hadoop.hdfs.server.blockmanagement.BlockPlacementPolicy and org.apache.hadoop.net.NetworkTopology 2022-08-25 10:21:38,083 WARN org.apache.hadoop.hdfs.protocol.BlockStoragePolicy: Failed to place enough replicas: expected size is 1 but only 0 storage types can be selected (replication=3, selected=[], unavailable=[DISK], removed=[DISK], policy=BlockStoragePolicy{HOT:7, storageTypes=[DISK], creationFallbacks=[], replicationFallbacks=[ARCHIVE]}) 2022-08-25 10:21:38,083 WARN org.apache.hadoop.hdfs.server.blockmanagement.BlockPlacementPolicy: Failed to place enough replicas, still in need of 1 to reach 3 (unavailableStorages=[DISK], storagePolicy=BlockStoragePolicy{HOT:7, storageTypes=[DISK], creationFallbacks=[], replicationFallbacks=[ARCHIVE]}, newBlock=true) All required storage types are unavailable: unavailableStorages=[DISK], storagePolicy=BlockStoragePolicy{HOT:7, storageTypes=[DISK], creationFallbacks=[], replicationFallbacks=[ARCHIVE]} 2022-08-25 10:21:38,083 INFO org.apache.hadoop.hdfs.StateChange: BLOCK* allocate blk_1073743799_2975, replicas=10.10.10.4:9866, 10.10.10.3:9866 for /benchmarks/TestDFSIO/io_control/in_file_test_io_1974 2022-08-25 10:21:38,085 INFO org.apache.hadoop.hdfs.StateChange: DIR* completeFile: /benchmarks/TestDFSIO/io_control/in_file_test_io_1974 is closed by DFSClient_NONMAPREDUCE_-1180995040_1
-
大佬们好,我们再对接华为大数据平台【FusionInsight Manager】时出现了一下问题问题描述:我们设计的Yarn任务提交设计以下几个步骤:检测 Yarn执行资源是否充足 【成功】QueueInfo queueInfo = yarnClient.getQueueInfo(amClientContext.getQueueName());设置yarn运行相关信息【成功】//部分代码 appContext.setApplicationName(amClientContext.getAppName()); appContext.setAttemptFailuresValidityInterval(20000); Set tags = new HashSet<>(1); tags.add("ddmp"); appContext.setApplicationTags(tags); ApplicationId appId = appContext.getApplicationId();上传待运行的任务至HDFS 【成功】 以下是部分代码,上传资源,包括设置yarn执行相关的环境变量,将AppMaster任务信息设置好/** * 添加一个本地资源到远程 * * @param fs 文件系统 * @param fileSrcPath 要上传的文件 * @param fileName 文件名 * @param appId 应用id * @param localResources 本地文件资源映射 * @param resources 文件资源 ,有时候我们并没有实际的资源信息,只有一个类似于命令操作,如果我们想将该命令生成一个文件并上传,就可以将该命令写在这里 * @throws IOException 异常信息 */ private void addToLocalResources(String appName, FileSystem fs, String fileSrcPath, String fileName, String appId, Map localResources, String resources) throws IOException { //获取要上传的目录路径 String suffix = appName + "/" + appId + "/" + fileName; Path dst = new Path(fs.getHomeDirectory(), suffix); //当要上传的文件不存在的时候 尝试将 resources 文件写入到一个目录中 if (fileSrcPath == null) { FSDataOutputStream ostream = null; try { //赋予 可读,可写,可执行的权限 ostream = FileSystem.create(fs, dst, new FsPermission((short) 456)); ostream.writeUTF(resources); } finally { IOUtils.closeStream(ostream); } } else { //将要上传的文件拷贝到对应的目录中 fs.copyFromLocalFile(new Path(fileSrcPath), dst); } //获取刚刚上传的文件的状态 FileStatus scFileStatus = fs.getFileStatus(dst); //创建一个本地资源映射 hdfs URI uri = dst.toUri(); URL url = URL.fromURI(uri); long len = scFileStatus.getLen(); long modificationTime = scFileStatus.getModificationTime(); LocalResource scRsrc = LocalResource.newInstance(url, LocalResourceType.FILE, LocalResourceVisibility.APPLICATION, len, modificationTime); //放入到资源映射中 localResources.put(fileName, scRsrc); }提交AppMaster任务到Yarn引擎 【失败】// 为应用程序主机设置容器启动上下文 ContainerLaunchContext amContainer = ContainerLaunchContext.newInstance(localResourceMap, env, commands, null, null, null); //权限处理 securityCheck(amContainer, amClientContext); //将容器设置进上下文对象 appContext.setAMContainerSpec(amContainer); //配置任务优先级状态 Priority pri = Priority.newInstance(0); appContext.setPriority(pri); //配置队列名称 appContext.setQueue(amClientContext.getQueueName()); yarnRunCallHook.doMessage("任务准备完成,开始提交任务!"); yarnClient.submitApplication(appContext);程序再运行到 yarnClient.submitApplication(appContext); 时执行卡住,通过日志观察,出现一下日志:48833 [main] INFO org.apache.hadoop.io.retry.RetryInvocationHandler - com.google.protobuf.InvalidProtocolBufferException: Protocol message end-group tag did not match expected tag., while invoking ApplicationClientProtocolPBClientImpl.getApplicationReport over 27. Trying to failover immediately. 48833 [main] INFO org.apache.hadoop.yarn.client.ConfiguredRMFailoverProxyProvider - Failing over to 28 49849 [main] INFO org.apache.hadoop.io.retry.RetryInvocationHandler - java.net.ConnectException: Call From DESKTOP-BTSFCSH/10.0.55.152 to 10-0-120-162:26004 failed on connection exception: java.net.ConnectException: Connection refused: no further information; For more details see: http://wiki.apache.org/hadoop/ConnectionRefused, while invoking ApplicationClientProtocolPBClientImpl.getApplicationReport over 28 after 1 failover attempts. Trying to failover after sleeping for 35465ms. 85315 [main] INFO org.apache.hadoop.yarn.client.ConfiguredRMFailoverProxyProvider - Failing over to 27 85366 [main] INFO org.apache.hadoop.io.retry.RetryInvocationHandler - com.google.protobuf.InvalidProtocolBufferException: Protocol message end-group tag did not match expected tag., while invoking ApplicationClientProtocolPBClientImpl.getApplicationReport over 27 after 2 failover attempts. Trying to failover after sleeping for 30581ms.请重点关注 Protocol message end-group tag did not match expected tag. 连接主节点的时候,出现协议不一致的问题连接信息如下:fs.defaultFS=hdfs://hacluster yarn.resourcemanager.address.27=10-0-120-161:26004 yarn.resourcemanager.address.28=10-0-120-162:26004 yarn.resourcemanager.ha.rm-ids=27,28 dfs.client.failover.proxy.provider.hacluster=org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider yarn.resourcemanager.scheduler.address.28=10-0-120-162:26002 dfs.nameservices=hacluster yarn.resourcemanager.scheduler.address.27=10-0-120-161:26002 dfs.namenode.rpc-address.hacluster.14=10-0-120-161:25000 dfs.namenode.rpc-address.hacluster.15=10-0-120-162:25000 yarn.resourcemanager.ha.enabled=true yarn.resourcemanager.recovery.enabled=true yarn.log-aggregation-enable=true dfs.ha.namenodes.hacluster=14,15 yarn.http.policy=HTTPS_ONLYFusionInsight Manager 已经开启Kereros,再本次提交中,kerberos认证已经通过 以上配置信息来自于 FusionInsight Manager 配置,确认端口信息等无误!以下是引入的Maven依赖 3.1.1 1.3.1 3.1.0 8 8 org.apache.hadoop hadoop-common ${hadoop.version} org.apache.hadoop hadoop-client ${hadoop.version} org.apache.hadoop hadoop-mapreduce-client-app ${hadoop.version} org.apache.hadoop hadoop-mapreduce-client-common ${hadoop.version} org.apache.hadoop hadoop-mapreduce-client-core ${hadoop.version} org.apache.hbase hbase-client ${hbase.version} org.apache.hbase hbase-common ${hbase.version} org.apache.hbase hbase-protocol ${hbase.version} org.apache.hbase hbase-server ${hbase.version} org.apache.hive hive-jdbc ${hive.version} org.apache.hive hive-service ${hive.version} 上述依赖,模仿华为云大数据平台 客户端案例的依赖!
-
伴随着科技的飞速发展,人工智能逐渐进入日常生活的各个方面。而大数据技术的研究和发展,则更推动技术的革新和社会经济的变革。大数据技术的出现背景、发展历程、研究现状以及发展过程中的存在问题是什么?同时在人工智能领域的大数据技术的发展又有哪些应用场景?让我们一起去探索。大数据的起源和发展随着互联网的广泛运用,云计算时代已经逐渐步入人们的生活,大数据在此背景下应运而生。1982年,约翰·奈斯比特在其著作中提出“我们现在大量生产信息,正如过去我们大量生产汽车一样”;阿尔文·托夫勒在《第三次浪潮》一书中,称大数据为“第三次浪潮的华彩乐章”;面对海量的数据,原有的处理方式已无法应对。2011年,麦肯锡全球研究所发布了《大数据:创新、竞争和生产力的下一个前沿》的报告,对“大数据”进行清晰解释;2012年,瑞士达沃斯召开世界经济论坛,大数据是会议主题。大数据发展起始于18世纪80年代初至90年代末,统计学家赫尔曼做出一台电动设备来统计美国本土人口普查数据,揭开数据处理新时代。雷德和普赖斯分别在1944年和1961年出版了《学者与研究型图书馆的未来》和《巴比伦以来的科学》,预测大数据时代的到来。2001年,美国Cartner公司推出大数据模型。2008年,美国自然杂志出版的一期专刊中第一次提出大数据——Big Data模型。大数据技术研究现状在国外,大数据技术被认为源于谷歌,在2003至2006年先后公开发表关于MapReduce、GFS和BigTable等核心技术学术论文。2012年,美国白宫颁布了《大数据研究与发展计划》,投入巨资到大数据研究领域。美国防部还开展XDATA项目,将大数据研究投入军事领域数据分析。在国内,2013年被称为大数据元年。2014年,国内众多互联网企业如小米、百度、腾讯、阿里等已将大数据技术应用于公司业务。大数据技术存在问题大数据技术的发展对各行各业也有重大影响,同时大数据技术的研究目前还不够完善,也面临着诸多问题需要去解决。(1)数据分析和处理问题。传统数据处理方式适用于少量的、结构化数据,而生活中采集的大多数数据是非结构化的数据,这就对数据分析和处理过程产生很大的影响。MapReduce计算并不能解决大数据处理的所有问题,需要更深层次的研究,解决数据处理的局限性。(2)数据安全和隐私问题。伴随着海量数据的采集和处理,数据的安全和数据的隐私问题应运而生。如何能够保障被采集数据的安全性、数据本身的隐私保护,都将是今后大数据技术研究的重点问题。现今社会就存在许多用户数据信息被盗用或者共享公开化等情况,这样对用户的人身和财产安全都产生很大的威胁。(3)政策和法规保障问题。面对大数据飞速发展的今天,国家缺乏相应监管体系,致使大数据滥用后用户的个人权益无法得到有效保障,国家应该建立相应的法律体系,对大数据的收集、开发和利用进行严格管理,同时对数据的正确使用设置规范和标准,推动大数据发展的规范化、合法化。大数据技术在人工智能领域的应用大数据技术在人工智能领域的应用广泛,涉及智慧农业、智慧城市、智慧工业等诸多方面[。智慧农业,大数据技术结合人工智能技术,收集海量数据信息进行处理分析,建立起精准农业、农产品流通体系、农业气象预测、农业环境管理等多个系统,推动农业的生产。智慧农业的提出使得很多技术可以整合使用,通过对土壤数据信息的采集,对土地耕作环境的监控,密切关注着温度、湿度的变化,并及时反馈监测的数据信息,预测出今后发展的方向。智慧城市,大数据技术应用于城市建设,建立城市数据信息共享平台,实时监管交通状况系统、智慧社区服务平台、城市地下排水监控系统等,推动城市智能化管理。智慧城市的格局要以网络化覆盖为基础,涵盖公共、卫生、交通、社区服务、社会保障等诸多方面,每一个环节都有相应的系统进行建构网络,然后通过各环节的相互配合去建造智慧城市格局。智慧工业,大数据技术促进工业“跨尺度、产业链、跨界”多源数据融合特点,推动工业化的智能管理和产量提升。工业化的进程在大数据技术的协助下,可以有效进行数据规整,对工业流程数据进行密切监控,确保产品生产环节精确度更高。当今时代是大数据的时代,大数据的合理使用将推动生产、生活的方方面面。大数据的发展和革新还在不断地发生改变,与此同时对大数据处理技术的研究也一直从未间断。然而,目前大数据的研究还处于初级阶段,很多技术不够完善,也存在着诸多问题,面临着巨大的挑战。伴随着智能化时代的来临,人工智能与大数据技术的结合将是今后大数据发展研究的重要主题。
上滑加载中
推荐直播
-
华为云码道Agent集成与鸿蒙实战2026/08/11 周二 19:00-21:00
王一男-华为云码道产品规划专家;李炎-华为云码道产品专家;彭江敏-华为云鸿蒙端云一体化开发专家
本次直播带你解读华为云码道7月份产品新特性、新功能。更有专家演示码道Agent Space × 钉钉机器集成实战,从0到1打通消息通道;码道鸿蒙端云一体化实战,快速搭建员工签到系统。
回顾中 -
华为云开发者AI素养直播课·第五期2026/09/04 周五 16:00-18:00
林华鼎-华为云AI开发者运营负责人;蒋春阳-华为云AI开发者案例开发专家
本期直播内容: AI工具体验营 · 第5-8课连讲。Agent-Team 多智能体协作完成毕业设计实践
回顾中 -
华为云开发者AI素养ClassRoom·第六期2026/09/08 周二 19:00-20:00
樊渊-2026华为软件挑战赛冠军
高手来了:看软挑高手解析二维排样问题—从工业难题到算法突破
回顾中
热门标签