首页 文章 精选 留言 我的

精选列表

搜索[实践],共10001篇文章
优秀的个人博客,低调大师

Android开发实践小结

作为一名搬运工,应该懂得避免重复创建轮子。 配置keystore密码信息 通常在app/build.gradle中我们会使用以下方式配置: signingConfigs { release { storeFile file("myapp.keystore") storePassword "mystorepassword" keyAlias "mykeyAlias" keyPassword "mykeypassword" } } 但这种方法不是特别好,因为如果你把代码分享到github,你的密码就泄露了。 推荐的做法应该是在Androd项目中gradle.properties(如果没有则手动创建一个)文件中创建以下变量,这个文件是不会被版本控制系统提交的,所以不用担心密码泄露。 KEYSTORE_PASSWORD=mystorepassword KEY_PASSWORD=mykeypassword 那么在app/build.gradle中的配置就应该是 signingConfigs { release { try { storeFile file("myapp.keystore") storePassword KEYSTORE_PASSWORD keyAlias "mykeyAlias" keyPassword KEY_PASSWORD } catch (ex) { throw new InvalidUserDataException("You should define KEYSTORE_PASSWORD and KEY_PASSWORD in gradle.properties.") } } } 这里在打包的时候,如果没有在gradle.properties配置相关变量,那么就会在messages窗口中提示 Error:(16, 0) You should define KEYSTORE_PASSWORD and KEY_PASSWORD in gradle.properties. 包引用尽可能使用dependency而不是直接导入jar包 dependencies { compile 'com.squareup.retrofit2:retrofit:2.1.0' } 以下方式也是不推荐使用的 dependencies { compile 'com.squareup.retrofit2:retrofit:2.1.0+' } 这种方式每次编译时都会去联网检测最新的包,影响性能,还可能会下载跟你不符合需求的jar包。 为release版本和debug版本指定不同的包名 同一个应用而不同的包名可以同时安装在同一个手机里面。 可以参考以下配置 android { buildTypes { debug { applicationIdSuffix '.debug' versionNameSuffix '-DEBUG' } ​ release { // ... } } } 分别为包名和版本号加上debug和DEBUG后缀。 如果想要在版本号添加时间信息,有利于区分,可以这样处理: 1、首先在app/build.gralde中定义一个buildTime()函数 //定义build 时间 def buildTime() { Date date = new Date() String build = date.format("yyMMddHHmm", TimeZone.getDefault()) return build } 2、然后再修改versionNameSuffix参数 buildTypes { debug { //配置这包名后缀下可以同时安装release包和debug包 applicationIdSuffix ".debug" //配置versionName后缀 versionNameSuffix '_' + (buildTime()) + '-DEBUG' ... } release { ... } } 在debug包中定义了版本号 开发调试工具 Stetho Stetho是facebook开源的Android调试工具,可以使用Chrome开发工具来对Android应用进行调试、抓包、查看Sqlite数据库等功能。可以在debug版本中集成Stetho,方便开发调试。 集成Stetho也是非常简单,只需要在app/build.gradle中配置 dependencies { compile 'com.facebook.stetho:stetho:1.4.1' } 然后在Application里面初始化Stetho public class MyApplication extends Application { public void onCreate() { super.onCreate(); Stetho.initializeWithDefaults(this); } } 这样就配置好了,AS连接手机跑起来后。打开Chrome,在地址栏输入 chrome://inspect/#devices 这时候就看到手机调试的信息 查看设备 查看sqlite数据库 View Hierarchy 调试网络 目前Steho支持okhttp网络库,同样的在gradle里面配置 dependencies { compile 'com.facebook.stetho:stetho:1.4.1' compile 'com.facebook.stetho:stetho-okhttp3:1.4.1' compile 'com.facebook.stetho:stetho-urlconnection:1.4.1' } 对OkHttp2.x版本 OkHttpClient client = new OkHttpClient(); client.networkInterceptors().add(new StethoInterceptor()); 对OkHttp3.x版本 new OkHttpClient.Builder() .addNetworkInterceptor(new StethoInterceptor()) .build(); 对于其他网络需要修改Stetho现在还没有支持,如果是HttpURLConnection,可以使用StethoURLConnectionManager来集成,详情可以参考Steho官网。 网络抓包 LeakCanary LeakCanary可以在应用运行时检测应用是否有OOM风险的一个工具库。同样的可只在debug版本中集成。 在app/build.gradle中 dependencies { debugCompile 'com.squareup.leakcanary:leakcanary-android:1.5' releaseCompile 'com.squareup.leakcanary:leakcanary-android-no-op:1.5' testCompile 'com.squareup.leakcanary:leakcanary-android-no-op:1.5' } 在Application中 public class ExampleApplication extends Application { ​ @Override public void onCreate() { super.onCreate(); if (LeakCanary.isInAnalyzerProcess(this)) { // This process is dedicated to LeakCanary for heap analysis. // You should not init your app in this process. return; } LeakCanary.install(this); // Normal app init code... } } 这样就配置好了。连接手机AS跑完后就会在手机桌面上生成一个LeakCanary图标,点击进入可以方便查看变量在内存中的引用链。如果LeakCanary检测到有内存泄露,也会发送一个通知栏消息来提醒。 AS常用插件 很多App都会使用UI注解框架来初始化UI控件其中最有名的估计就是ButterKnife了。 class ExampleActivity extends Activity { @BindView(R.id.user) EditText username; @BindView(R.id.pass) EditText password; ​ @BindString(R.string.login_error) String loginErrorMessage; ​ @OnClick(R.id.submit) void submit() { // TODO call server... } ​ @Override public void onCreate(Bundle savedInstanceState) { super.onCreate(savedInstanceState); setContentView(R.layout.simple_activity); ButterKnife.bind(this); // TODO Use fields... } } 但是经常写@BindView也是一件令人枯燥和烦恼的事情。 Android Butterknife Zelezny 这个插件可以极大的解放程序猿的双手,提高搬砖效率。 安装好插件后,把光标定位到layout文件的引用处,例如setContentView(R.layout.simple_activity);的R.layout.simple_activity末尾处,按下快捷键command+N(windows是Altt+Insert)将弹出以下对话框。 so easy! 关注我们,可以获取更多

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

Kylin-实践OLAP

OLAP的历史与基本概念 OLAP全称为在线联机分析应用,是一种对于多维数据分析查询的解决方案。典型的OLAP应用场景包括销售、市场、管理等商务报表,预算决算,经济报表等等。 最早的OLAP查询工具是发布于1970年的Express,然而完整的OLAP概念是在1993年由关系数据库之父EdgarF.Codd 提出,伴随而来的是著名的“twelvelaws of online analytical processing”. 1998年微软发布MicrosoftAnalysis Services,并且在早一年通过OLE DB for OLAP API引入MDX查询语言,2001年微软和Hyperion发布的XML forAnalysis 成为了事实上的OLAP查询标准。如今,MDX已成为与SQL旗鼓相当的OLAP 查询语言,被各家OLAP厂商先后支持。 OLAPCube是一种典型的多维数据分析技术,Cube本身可以认为是不同维度数据组成的dataset,一个OLAP Cube 可以拥有多个维度(Dimension),以及多个事实(Factor Measure)。用户通过OLAP工具从多个角度来进行数据的多维分析。通常认为OLAP包括三种基本的分析操作:上卷(rollup)、下钻(drilldown)、切片切块(slicingand dicing),原始数据经过聚合以及整理后变成一个或多个维度的视图。 ROLAP和MOLAP 传统OLAP根据数据存储方式的不同分为ROLAP(Relational OLAP)以及MOLAP(Multi-dimensionOLAP) ROLAP 以关系模型的方式存储用作多维分析用的数据,优点在于存储体积小,查询方式灵活,然而缺点也显而易见,每次查询都需要对数据进行聚合计算,为了改善短板,ROLAP使用了列存、并行查询、查询优化、位图索引等技术 MOLAP 将分析用的数据物理上存储为多维数组的形式,形成CUBE结构。维度的属性值映射成多维数组的下标或者下标范围,事实以多维数组的值存储在数组单元中,优势是查询快速,缺点是数据量不容易控制,可能会出现维度爆炸的问题。 大数据时代OLAP的挑战 近二十年内,ROLAP技术随着MPP并行数据库技术的发展,尤其是列存技术的支持下,实现了分析能力大幅度的跨越提升,同时伴随着内存成本的进一步降低,单节点内存扩展性增强,集群单节点的查询性能实现了飞跃,内存数据库的实用性跨上了一个新台阶,这些技术进步共同作用的结果是类似的技术基本覆盖了TB级别的数据分析需求。 Hadoop以及相关大数据技术的出现提供了一个几近无限扩展的数据平台,在相关技术的支持下,各个应用的数据已突破了传统OLAP所能支持的容量上界。每天千万、数亿条的数据,提供若干维度的分析模型,大数据OLAP最迫切所要解决的问题就是大量实时运算导致的响应时间迟滞。 2. Apache Kylin 大数据下的OLAP解决方案 Apache Kylin的背景 Apache Kylin 是一个Hadoop生态圈下的MOLAP系统,是eBay大数据部门从2014年开始研发并开源的支持TB到PB级别数据量的分布式OLAP分析引擎。其特点包括: 可扩展的超快的OLAP引擎 提供ANSI-SQL接口 交互式查询能力 MOLAP Cube 的概念 与BI工具可无缝整合 Apache Kylin典型的应用场景如下: 用户数据存在于Hadoop HDFS中,利用Hive将HDFS文件数据以关系数据方式存取,数据量巨大,在500G以上 每天有数G甚至数十G的数据增量导入 有10个左右为固定的分析维度 ApacheKylin的核心思想是利用空间换时间,由于查询方面制定了多种灵活的策略,进一步提高空间的利用率,使得这样的平衡策略在应用中是值得采用的。 Apache Kylin的总体架构 Apache Kylin 作为一个OLAP引擎完成了从数据源抓取数据,ETL到自己的存储引擎,提供REST服务等一系列工作,其架构如图所示: Apache Kylin 的生态圈包括: Kylin Core: Kylin 引擎的框架,查询、任务、以及存储引擎都集中于此,除此之外还包括一个REST 服务器来响应各种客户端请求。 扩展插件: 各种提供额外特性的插件,如安全认证、SSO等 完整性组件: Job管理器,ETL、监控以及报警 交互界面: 基于Kylin Core之上的用户交互界面 驱动: 提供了JDBC以及ODBC的连接方式 Apache kylin Cube 多维数据的计算 Apache Kylin的多维计算主要是体现在OLAPCube的计算。Cube由多个Cuboid组合而成,Cuboid上的数据是原始数据聚合的数据,因此创建Cube可以看作是在原始数据导入时做的一个预计算预处理的过程。Kylin的强大之处在于充分利用了Hadoop的MapReduce并行处理的能力,高效处理导入的数据。 Apache Kylin的数据来自于Hive,并作为一个Hive的加速器希望最终的查询SQL类似于直接在Hive上查询。因此Kylin在建立Cube的时候需要从Hive获取Hive表的元数据。虽然有建立Cube的过程,但是并不想对普通的查询用户暴露Cube的存在。 Apache Kylin创建Cube的过程如下图所示: 根据Cube定义的事实表以及维度表,利用Hive创建一张宽表 抽取事实表上的维度的distinct值,将事实表上的维度以字典树方式压缩编码成目录,将维度表以字典树的方式编码 利用MapReduce从第一步得到的宽表文件作为输入,创建 N-Dimension cuboid,然后每次根据前一步的结果串行生成 N-1 cuboid, N-2 cuboid … 0-Cuboid 根据生成的Cuboid数据量计算HTable的Region分割策略,创建HTable,将HFile导入进来 Apache Kylin与传统的OLAP一样,无法应对数据Update的情况(更新数据会导致Cube的失效,需要重建整个Cube)。面对每天甚至每两个小时这样固定周期的增量数据,Kylin使用了一种增量Cubing技术来进行快速响应。 Apache Kylin的Cube可以根据时间段划分成多个Segment。在Cube第一次Build完成之后会有一个Segment,在每次增量Build后会产生一个新的Segment。增量Cubing依赖已有的CubeSegments和增量的原始数据。增量Cubing的步骤和新建 Cube的步骤类似,Segment之间以时间段进行区分。 增量Cubing所需要面对的原始数据量更小,因此增量Cubing的速度是非常快的。然而随着CubeSegments的数目增加,一定程度上会影响到查询的进行,所以在Segments数目到一定数量后可能需要进行CubeSegments的合并操作,实际上MergeCube是合成了一个新的大的CubeSegment来替代,Merge操作是一个异步的在线操作,不会对前端的查询业务产生影响。 合并操作步骤如下: 遍历指定的Cube Segment 合并维度字典目录和维度表快照 利用MapReduce合并他们的 N-Dimension cuboid 将cuboid转换成HFile,生成新的HTable,替代原有的多个HTable Apache Kylin对传统MOLAP的改进 计算Cube的存储代价以及计算代价都是比较大的, 传统OLAP的维度爆炸的问题Kylin也一样会遇到。 Kylin提供给用户一些优化措施,在一定程度上能降低维度爆炸的问题: Cube 优化: Hierachy Dimension Derived Dimension Aggregation Group Hierachy Dimension, 一系列具有层次关系的Dimension组成一个Hierachy, 比如年、月、日组成了一个Hierachy, 在Cube中,如果不设置Hierarchy, 会有 年、月、日、年月、年日、月日 6个cuboid, 但是设置了Hierarchy之后Cuboid增加了一个约束,希望低Level的Dimension一定要伴随高Level的Dimension 一起出现。设置了Hierachy Dimension 能使得需要计算的维度组合减少一半。 Derived Dimension, 如果在某张维度表上有多个维度,那么可以将其设置为Derived Dimension, 在Kylin内部会将其统一用维度表的主键来替换,以此来达到降低维度组合的数目,当然在一定程度上Derived Dimension 会降低查询效率,在查询时,Kylin使用维度表主键进行聚合后,再通过主键和真正维度列的映射关系做一次转换,在Kylin内部再对结果集做一次聚合后返回给用户 Aggregation Group, 这是一个将维度进行分组,以求达到降低维度组合数目的手段。不同分组的维度之间组成的Cuboid数量会大大降低,维度组合从2的(k+m+n)次幂至多能降低到 2的k次幂加2的m次幂加2的n次幂。Group的优化措施与查询SQL紧密依赖,可以说是为了查询的定制优化。 如果查询的维度是夸Group的,那么Kylin需要以较大的代价从N-Cuboid中聚合得到所需要的查询结果,这需要Cube构建人员在建模时仔细地斟酌。 数据压缩: ApacheKylin针对维度字典以及维度表快照采用了特殊的压缩算法,对于Hbase中的聚合计算数据利用了Hadoop的LZO或者是Snappy,从而保证存储在Hbase以及内存中的数据尽可能的小。其中维度字典以及维度表快照的压缩考虑到DataCube中会出现非常多的重复的维度成员值,最直接的处理方式就是利用数据字典的方式将维度值映射成ID, Kylin中采用了Trie树的方式对维度值进行编码 distinct count聚合查询优化: Apache Kylin 采用了HypeLogLog的方式来计算DistinctCount。好处是速度快,缺点是结果是一个近似值,会有一定的误差。在非计费等通常的场景下DistinctCount的统计误差应用普遍可以接受。 具体的算法可见Paper,本文不再赘述: http://algo.inria.fr/flajolet/Publications/FlFuGaMe07.pdf Apache kylin SQL查询的实现 ANSI SQL查询是Apache Kylin 非常明显的优势。Kylin的SQL语法解析依赖于另一个开源数据管理框架 ApacheCalcite, Calcite即之前的Optiq,是一个没有存储模块的数据库,即不管理数据存储、不包含数据处理的算法,不包含元信息的存储。因此它非常适合来做一个应用到存储引擎之间的中间层。在Calcite的基础之上只要为存储引擎写一个专用的适配器(Adapter)即可形成一个功能丰富的支持DML甚至DDL的“类数据库”。 Kylin完成了一个定制的Adapter,在Calcite完成SQL解析,形成语法树(AST)之后,由Kylin定义语法树各个节点的执行规则来进行查询。Calcite在遍历语法树节点后生成一个Kylin描述查询模型的Digest, Kylin会为此Digest去判断是否有匹配的Cube。如果有与查询匹配的Cube,即选择一个查询代价最小的Cube进行查询(KylinCube的查询代价计算目前是一个开放接口,可以根据维度数目,可以根据数据量大小来计算Cost) Kylin目前的多维数据存储引擎是HBase, Kylin利用了HBase的Coprocessor机制在HBase的RegionServer完成部分聚合以及全部过滤操作,在HbaseScan时提前进行计算,利用HBase多个Region Server的计算能力加速Kylin的SQL查询。目前Kylin仍然有部分查询语法不支持,特别是过滤器Where部分的约束较多、对SQL有一定的要求,但是如果有针对性的对Coprocessor部分进行改造相信SQL兼容度可以有大幅的提升。 Apache kylin 与 RTOLAP ApacheKylin 可以说是与市面上流行的Presto、SparkSQL、Impala等直接在原始数据上查询的系统(暂且归于RTOLAP)走了一条完全不同的道路。前者在如何快速求得预计算结果,以及优化查询解析使得更多的查询能用上预计算结果方面在优化。后续Kylin的版本会改进预计算引擎,优化预计算速度,使得Kylin可以变成一个近似实时的分析引擎。而像Presto,SparkSQL等是着重于优化查询数据的过程环节,像一些其它的数据仓库一样,使用列存、压缩、并行查询等技术,优化查询。这种方案的好处就在于扩展性强、能适配更广泛的查询。但是在查询速度上,可以说Apache Kylin 要比ROLAP 至少快上一个数量级,所以对与查询响应时间要求较高的应用,ApacheKylin是最好的选择。 3. Apache Kylin在网易 Kylin服务化 在网易,Apache Kylin作为大数据平台的OLAP查询模块,可以为公司的各种分析类需求以及应用提供服务。所有数据存在Hadoop Hive 上的数据都能够通过Kylin OLAP 引擎进行加速查询。在公司内部Kylin作为一个统一平台,与各产品的数据仓库进行接驳。 目前Kylin的部署架构如下: Kylin集群由多个查询节点以及控制节点组成。 控制节点唯一,负责集群项目、任务调度与Cube增删查改。 多个查询节点前用Nginx做负载均衡,后段节点可按需水平扩容。前端可同时支持JDBC与ODBC的客户端查询。 Kylin性能表现 在Kylin上线前,我们选取了公司内部原有的一些报表业务进行过性能对比,对比内容在相同的数据下、Kylin查询与Mondrian 结合Oracle的查询比较。 测试结果通过数据量较大的DataStream报表来进行比较: 再看Kylin的吞吐量,利用Haproxy进行请求转发后随着Kylin服务器的增加吞吐量的表现: 网易对Kylin的改进 原生的社区版Aapche Kylin 是需要部署在一个统一底层的Hadoop、Hive、HBase集群之上的。而网易内部的大数据平台由于各种原因,分为了多个Hadoop集群、各应用会在不同的Hadoop集群上建立Hive数据仓库。最原始而自然的想法就是在每一个Hadoop环境上部署一套Kylin服务来满足不同的需求,但是集群资源管理、计算资源调度、管理运维的复杂性都会是一个比较突出的问题。例如用户数据在A机房的Hive上,而A机房的Hadoop集群并没有足够的计算资源来保证KylinOLAP的高效运行。因此根据公司内部实际的大数据平台分布情况及机房建设情况,将Kylin打造成一个公司内统一的服务平台是一个更好的选择。OLAP小组对开源版本的Kylin进行了二次开发,并将改进补丁提交给了社区并受到了积极反馈。 目前的改进主要包括: Kylin对Kerberos认证的支持 Kylin非Hadoop节点的部署支持 多数据源的支持 在公司内,由于性能以及安全性方面的考量,不同部门的应用会搭建各自的Hive进行数据分析,并且由于公司内还没有跨机房的Hadoop集群,因此会出现用户数据在A地方的Hive上,而A机房的Hadoop集群并没有足够的计算资源来保证KylinOLAP的高效运行。 综合分析现实的场景之后,我们选择了公司内最大的hadoop集群作为KylinOLAP的计算引擎集群,保证有充足的存储以及计算资源。 HBase采用一个独立的集群,避免Hbase查询和Hadoop集群任务之间的互相干扰。数据源Hive允许用户自定义,目前已支持同Hadoop集群下不同Hive 以及不同Hadoop集群下的不同Hive节点使用KylinOLAP服务。根据用户数据仓库的实际配置情况可能会出现跨集群的数据源抽取计算,由于公司同城机房有专线网络,数据仓库Hive里的源数据量也远小于Kylin实际的聚合后的数据存储(存于Hbase,数据量大小一般为数据源Hive中的10倍以上),因此可认为这样的开销可以认为带来的影响不大,并且在我们的测试中得到了印证。 Kylin OLAP与猛犸以及有数的结合 猛犸是网易内部的统一大数据入口平台,为了让Kylin更快更好的融入到大平台中,OLAP小组已计划在不久之后全面与猛犸大数据平台进行打通和整合, Kylin OLAP将深度内嵌于猛犸,用户可以基于猛犸平台完成KylinOLAP的简化管理工作。猛犸平台对接控制节点,作为专业数据建模师的操作入口 Kylin将利用猛犸的用户管理功能 猛犸将接管用户项目的创建以及Cube的管理 猛犸将原有的Hive数据源彻底与Kylin打通,便于Kylin管理用户的数据源 Kylin原生的用户管理是基于LDAP的,如果不使用LDAP服务需要利用SpringSecurity重新开发一套,网易的内部的猛犸大数据平台有一套成熟且完善的用户权限访问控制体系,因此可以利用现成的机制对Kylin的访问、修改做保护性的限制。 Kylin的Data Cube建模,特别是一些高级的Cube优化功能如RowKey顺序、维度分组、分层等需要较高的学习成本,所以认为不适合让一般的数据分析师来直接操作,我们设计了一套简化版的Cube 建模流程,以用户申请——运维审批的方式进行数据的接入。 有数是网易内部重要的报表分析平台,有数将KylinOLAP作为一个单独的数据源进行支持。已有的以及潜在的Hive查询客户可以轻松的将报表迁移到KylinOLAP,使得大数据量下的交互式报表分析成为可能。 有数能基于在猛犸上创建的Cube创建报表 有数主动识别Kylin Cube定义的维度和度量 用户在Kylin OLAP允许的范围内自由操作,完成报表的编辑和查询。 与有数结合后的Kylin 查询结果可以用更多更丰富的图表的方式展示给数据分析人员:

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

Spark实践-日志查询

环境 win 7 jdk 1.7.0_79 (Oracle Corporation) scala version 2.10.5 spark 1.6.1 详细配置: Spark Properties spark.app.id local-1461891171126 spark.app.name JavaLogQuery spark.driver.host 10.170.26.123 spark.driver.port 34998 spark.executor.id driver spark.externalBlockStore.folderName spark-5242ec5b-3653-42e4-9ba2-da3ef515a1d5 spark.master local[1] spark.scheduler.mode FIFO 任务 完成对如下日志的查询: "10.10.10.10 - \"FRED\" [18/Jan/2013:17:56:07 +1100] \"GET http://images.com/2013/Generic.jpg " + "HTTP/1.1\" 304 315 \"http://referall.com/\" \"Mozilla/4.0 (compatible; MSIE 7.0; " + "Windows NT 5.1; GTB7.4; .NET CLR 2.0.50727; .NET CLR 3.0.04506.30; .NET CLR 3.0.04506.648; " + ".NET CLR 3.5.21022; .NET CLR 3.0.4506.2152; .NET CLR 1.0.3705; .NET CLR 1.1.4322; .NET CLR " + "3.5.30729; Release=ARP)\" \"UD-1\" - \"image/jpeg\" \"whatever\" 0.350 \"-\" - \"\" 265 923 934 \"\" " + "62.24.11.25 images.com 1358492167 - Whatup", "10.10.10.10 - \"FRED\" [18/Jan/2013:18:02:37 +1100] \"GET http://images.com/2013/Generic.jpg " + "HTTP/1.1\" 304 306 \"http:/referall.com\" \"Mozilla/4.0 (compatible; MSIE 7.0; Windows NT 5.1; " + "GTB7.4; .NET CLR 2.0.50727; .NET CLR 3.0.04506.30; .NET CLR 3.0.04506.648; .NET CLR " + "3.5.21022; .NET CLR 3.0.4506.2152; .NET CLR 1.0.3705; .NET CLR 1.1.4322; .NET CLR " + "3.5.30729; Release=ARP)\" \"UD-1\" - \"image/jpeg\" \"whatever\" 0.352 \"-\" - \"\" 256 977 988 \"\" " + "0 73.23.2.15 images.com 1358492557 - Whatup" 思路: 1.利用正则表达式提取出日志特征,然后map在分片后的RDD上。 JavaPairRDD<Tuple3<String, String, String>, Stats> extracted 2.执行reducebykey,merge相同的Stats package org.apache.spark.examples; import com.google.common.collect.Lists; import scala.Tuple2; import scala.Tuple3; import org.apache.commons.logging.impl.Log4JLogger; import org.apache.log4j.Level; import org.apache.log4j.Logger; import org.apache.spark.SparkConf; import org.apache.spark.api.java.JavaPairRDD; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.JavaSparkContext; import org.apache.spark.api.java.function.Function2; import org.apache.spark.api.java.function.PairFunction; import java.io.Serializable; import java.util.Collections; import java.util.List; import java.util.regex.Matcher; import java.util.regex.Pattern; /** * 日志查询 * @author jinhang * */ public final class JavaLogQuery { //模拟日志 exampleApacheLogs public static final List<String> exampleApacheLogs = Lists.newArrayList( "10.10.10.10 - \"FRED\" [18/Jan/2013:17:56:07 +1100] \"GET http://images.com/2013/Generic.jpg " + "HTTP/1.1\" 304 315 \"http://referall.com/\" \"Mozilla/4.0 (compatible; MSIE 7.0; " + "Windows NT 5.1; GTB7.4; .NET CLR 2.0.50727; .NET CLR 3.0.04506.30; .NET CLR 3.0.04506.648; " + ".NET CLR 3.5.21022; .NET CLR 3.0.4506.2152; .NET CLR 1.0.3705; .NET CLR 1.1.4322; .NET CLR " + "3.5.30729; Release=ARP)\" \"UD-1\" - \"image/jpeg\" \"whatever\" 0.350 \"-\" - \"\" 265 923 934 \"\" " + "62.24.11.25 images.com 1358492167 - Whatup", "10.10.10.10 - \"FRED\" [18/Jan/2013:18:02:37 +1100] \"GET http://images.com/2013/Generic.jpg " + "HTTP/1.1\" 304 306 \"http:/referall.com\" \"Mozilla/4.0 (compatible; MSIE 7.0; Windows NT 5.1; " + "GTB7.4; .NET CLR 2.0.50727; .NET CLR 3.0.04506.30; .NET CLR 3.0.04506.648; .NET CLR " + "3.5.21022; .NET CLR 3.0.4506.2152; .NET CLR 1.0.3705; .NET CLR 1.1.4322; .NET CLR " + "3.5.30729; Release=ARP)\" \"UD-1\" - \"image/jpeg\" \"whatever\" 0.352 \"-\" - \"\" 256 977 988 \"\" " + "0 73.23.2.15 images.com 1358492557 - Whatup"); public static final Pattern apacheLogRegex = Pattern.compile( "^([\\d.]+) (\\S+) (\\S+) \\[([\\w\\d:/]+\\s[+\\-]\\d{4})\\] \"(.+?)\" (\\d{3}) ([\\d\\-]+) \"([^\"]+)\" \"([^\"]+)\".*"); public static class Stats implements Serializable { private final int count; private final int numBytes; public Stats(int count, int numBytes) { this.count = count; this.numBytes = numBytes; } public Stats merge(Stats other) { return new Stats(count + other.count, numBytes + other.numBytes); } public String toString() { return String.format("bytes=%s\tn=%s", numBytes, count); } } public static Tuple3<String, String, String> extractKey(String line) { Matcher m = apacheLogRegex.matcher(line); if (m.find()) { String ip = m.group(1); String user = m.group(3); String query = m.group(5); if (!user.equalsIgnoreCase("-")) { return new Tuple3<String, String, String>(ip, user, query); } } return new Tuple3<String, String, String>(null, null, null); } public static Stats extractStats(String line) { Matcher m = apacheLogRegex.matcher(line); if (m.find()) { int bytes = Integer.parseInt(m.group(7)); return new Stats(1, bytes); } else { return new Stats(1, 0); } } public static void main(String[] args) { Logger.getLogger(JavaLogQuery.class).setLevel(Level.FATAL); SparkConf sparkConf = new SparkConf().setAppName("JavaLogQuery").setMaster("local[1]"); JavaSparkContext jsc = new JavaSparkContext(sparkConf); JavaRDD<String> dataSet = (args.length == 1) ? jsc.textFile(args[0]) : jsc.parallelize(exampleApacheLogs); JavaPairRDD<Tuple3<String, String, String>, Stats> extracted = dataSet.mapToPair(new PairFunction<String, Tuple3<String, String, String>, Stats>() { @Override public Tuple2<Tuple3<String, String, String>, Stats> call(String s) { return new Tuple2<Tuple3<String, String, String>, Stats>(extractKey(s), extractStats(s)); } }); JavaPairRDD<Tuple3<String, String, String>, Stats> counts = extracted.reduceByKey(new Function2<Stats, Stats, Stats>() { @Override public Stats call(Stats stats, Stats stats2) { return stats.merge(stats2); } }); List<Tuple2<Tuple3<String, String, String>, Stats>> output = counts.collect(); //遍历结果 for (Tuple2<?,?> t : output) { System.out.println(t._1() + "\t" + t._2()); } jsc.stop(); } } 分析下执行过程: 加载SLF4J Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties SLF4J: Class path contains multiple SLF4J bindings. SLF4J: Found binding in [jar:file:/D:/JavaProject/spark-demo/lib/spark-assembly-1.6.1-hadoop2.0.0-mr1-cdh4.2.0.jar!/org/slf4j/impl/StaticLoggerBinder.class] SLF4J: Found binding in [jar:file:/D:/JavaProject/spark-demo/lib/spark-examples-1.6.1-hadoop2.0.0-mr1-cdh4.2.0.jar!/org/slf4j/impl/StaticLoggerBinder.class] SLF4J: See http://www.slf4j.org/codes.html#multiple_bindings for an explanation. SLF4J: Actual binding is of type [org.slf4j.impl.Log4jLoggerFactory] 初始化sparkcontext上下文 16/04/29 09:36:22 INFO SparkContext: Running Spark version 1.6.1 //-Djava.library.path=$HADOOP_HOME/lib/native/Linux-amd64-64/*.jar可以解决 16/04/29 09:36:23 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable 16/04/29 09:36:23 INFO SecurityManager: Changing view acls to: hp 16/04/29 09:36:23 INFO SecurityManager: Changing modify acls to: hp 16/04/29 09:36:23 INFO SecurityManager: SecurityManager: authentication disabled; ui acls disabled; users with view permissions: Set(hp); users with modify permissions: Set(hp) 16/04/29 09:36:23 INFO Utils: Successfully started service 'sparkDriver' on port 36010. 16/04/29 09:36:23 INFO Slf4jLogger: Slf4jLogger started 16/04/29 09:36:24 INFO Remoting: Starting remoting 16/04/29 09:36:24 INFO Remoting: Remoting started; listening on addresses :[akka.tcp://sparkDriverActorSystem@10.170.26.123:36023] 16/04/29 09:36:24 INFO Utils: Successfully started service 'sparkDriverActorSystem' on port 36023. 16/04/29 09:36:24 INFO SparkEnv: Registering MapOutputTracker 16/04/29 09:36:24 INFO SparkEnv: Registering BlockManagerMaster 16/04/29 09:36:24 INFO DiskBlockManager: Created local directory at C:\Users\hp\AppData\Local\Temp\blockmgr-84667505-0018-439b-9627-a4360d872118 16/04/29 09:36:24 INFO MemoryStore: MemoryStore started with capacity 517.4 MB 16/04/29 09:36:24 INFO SparkEnv: Registering OutputCommitCoordinator 16/04/29 09:36:24 INFO Utils: Successfully started service 'SparkUI' on port 4040. 16/04/29 09:36:24 INFO SparkUI: Started SparkUI at http://10.170.26.123:4040 16/04/29 09:36:24 INFO Executor: Starting executor ID driver on host localhost 16/04/29 09:36:24 INFO Utils: Successfully started service 'org.apache.spark.network.netty.NettyBlockTransferService' on port 36030. 16/04/29 09:36:24 INFO NettyBlockTransferService: Server created on 36030 16/04/29 09:36:24 INFO BlockManagerMaster: Trying to register BlockManager 16/04/29 09:36:24 INFO BlockManagerMasterEndpoint: Registering block manager localhost:36030 with 517.4 MB RAM, BlockManagerId(driver, localhost, 36030) 16/04/29 09:36:24 INFO BlockManagerMaster: Registered BlockManager SecurityManager ‘sparkDriver’ on port 36010 Remoting: Remoting started; listening on addresses :[akka.tcp://sparkDriverActorSystem@10.170.26.123:36023] MapOutputTracker BlockManagerMaster DiskBlockManager: Created local directory at C:\Users\hp\AppData\Local\Temp\blockmgr-84667505-0018-439b-9627- OutputCommitCoordinator Executor org.apache.spark.network.netty.NettyBlockTransferService 这几个是几个主要过程。 开始执行job 16/04/29 10:12:31 INFO SparkContext: Starting job: collect at JavaLogQuery.java:112 16/04/29 10:12:31 INFO DAGScheduler: Registering RDD 1 (mapToPair at JavaLogQuery.java:98) 16/04/29 10:12:31 INFO DAGScheduler: Got job 0 (collect at JavaLogQuery.java:112) with 1 output partitions 16/04/29 10:12:31 INFO DAGScheduler: Final stage: ResultStage 1 (collect at JavaLogQuery.java:112) 16/04/29 10:12:31 INFO DAGScheduler: Parents of final stage: List(ShuffleMapStage 0) 16/04/29 10:12:31 INFO DAGScheduler: Missing parents: List(ShuffleMapStage 0) 16/04/29 10:12:31 INFO DAGScheduler: Submitting ShuffleMapStage 0 (MapPartitionsRDD[1] at mapToPair at JavaLogQuery.java:98), which has no missing parents 16/04/29 10:12:31 WARN SizeEstimator: Failed to check whether UseCompressedOops is set; assuming yes 16/04/29 10:12:31 INFO MemoryStore: Block broadcast_0 stored as values in memory (estimated size 3.1 KB, free 3.1 KB) 16/04/29 10:12:31 INFO MemoryStore: Block broadcast_0_piece0 stored as bytes in memory (estimated size 1897.0 B, free 5.0 KB) 16/04/29 10:12:31 INFO BlockManagerInfo: Added broadcast_0_piece0 in memory on localhost:36394 (size: 1897.0 B, free: 517.4 MB) 16/04/29 10:12:31 INFO SparkContext: Created broadcast 0 from broadcast at DAGScheduler.scala:1006 16/04/29 10:12:31 INFO DAGScheduler: Submitting 1 missing tasks from ShuffleMapStage 0 (MapPartitionsRDD[1] at mapToPair at JavaLogQuery.java:98) 16/04/29 10:12:31 INFO TaskSchedulerImpl: Adding task set 0.0 with 1 tasks 16/04/29 10:12:31 INFO TaskSetManager: Starting task 0.0 in stage 0.0 (TID 0, localhost, partition 0,PROCESS_LOCAL, 3033 bytes) 16/04/29 10:12:31 INFO Executor: Running task 0.0 in stage 0.0 (TID 0) 16/04/29 10:12:32 INFO Executor: Finished task 0.0 in stage 0.0 (TID 0). 1158 bytes result sent to driver 16/04/29 10:12:32 INFO TaskSetManager: Finished task 0.0 in stage 0.0 (TID 0) in 320 ms on localhost (1/1) 16/04/29 10:12:32 INFO TaskSchedulerImpl: Removed TaskSet 0.0, whose tasks have all completed, from pool 16/04/29 10:12:32 INFO DAGScheduler: ShuffleMapStage 0 (mapToPair at JavaLogQuery.java:98) finished in 0.349 s 16/04/29 10:12:32 INFO DAGScheduler: looking for newly runnable stages 16/04/29 10:12:32 INFO DAGScheduler: running: Set() 16/04/29 10:12:32 INFO DAGScheduler: waiting: Set(ResultStage 1) 16/04/29 10:12:32 INFO DAGScheduler: failed: Set() 16/04/29 10:12:32 INFO DAGScheduler: Submitting ResultStage 1 ***(ShuffledRDD[2] at reduceByKey at JavaLogQuery.java:105), which has no missing parents*** 16/04/29 10:12:32 INFO MemoryStore: Block broadcast_1 stored as values in memory (estimated size 2.9 KB, free 7.9 KB) 16/04/29 10:12:32 INFO MemoryStore: Block broadcast_1_piece0 stored as bytes in memory (estimated size 1746.0 B, free 9.6 KB) 16/04/29 10:12:32 INFO BlockManagerInfo: Added broadcast_1_piece0 in memory on localhost:36394 (size: 1746.0 B, free: 517.4 MB) 16/04/29 10:12:32 INFO SparkContext: Created broadcast 1 from broadcast at DAGScheduler.scala:1006 16/04/29 10:12:32 INFO DAGScheduler: Submitting 1 missing tasks from ResultStage 1 (ShuffledRDD[2] at reduceByKey at JavaLogQuery.java:105) 16/04/29 10:12:32 INFO TaskSchedulerImpl: Adding task set 1.0 with 1 tasks 16/04/29 10:12:32 INFO TaskSetManager: Starting task 0.0 in stage 1.0 (TID 1, localhost, partition 0,NODE_LOCAL, 1894 bytes) 16/04/29 10:12:32 INFO Executor: Running task 0.0 in stage 1.0 (TID 1) 16/04/29 10:12:32 INFO ShuffleBlockFetcherIterator: Getting 1 non-empty blocks out of 1 blocks 16/04/29 10:12:32 INFO ShuffleBlockFetcherIterator: Started 0 remote fetches in 16 ms 16/04/29 10:12:32 INFO Executor: Finished task 0.0 in stage 1.0 (TID 1). 1449 bytes result sent to driver 16/04/29 10:12:32 INFO TaskSetManager: Finished task 0.0 in stage 1.0 (TID 1) in 107 ms on localhost (1/1) 16/04/29 10:12:32 INFO TaskSchedulerImpl: Removed TaskSet 1.0, whose tasks have all completed, from pool 16/04/29 10:12:32 INFO DAGScheduler: ResultStage 1 (collect at JavaLogQuery.java:112) finished in 0.108 s 16/04/29 10:12:32 INFO DAGScheduler: Job 0 finished: collect at JavaLogQuery.java:112, took 0.850227 s (10.10.10.10,"FRED",GET http://images.com/2013/Generic.jpg HTTP/1.1) bytes=621 n=2 16/04/29 10:12:56 INFO BlockManagerInfo: Removed broadcast_1_piece0 on localhost:36394 in memory (size: 1746.0 B, free: 517.4 MB) 结束 16/04/29 10:12:56 INFO BlockManagerInfo: Removed broadcast_1_piece0 on localhost:36394 in memory (size: 1746.0 B, free: 517.4 MB) 16/04/29 10:16:13 INFO ContextCleaner: Cleaned accumulator 2 16/04/29 10:16:13 INFO BlockManagerInfo: Removed broadcast_0_piece0 on localhost:36394 in memory (size: 1897.0 B, free: 517.4 MB) 16/04/29 10:16:13 INFO ContextCleaner: Cleaned accumulator 1 16/04/29 10:24:29 WARN QueuedThreadPool: 5 threads could not be stopped 16/04/29 10:24:29 INFO SparkUI: Stopped Spark web UI at http://10.170.26.123:4040 16/04/29 10:24:29 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped! 16/04/29 10:24:29 INFO MemoryStore: MemoryStore cleared 16/04/29 10:24:29 INFO BlockManager: BlockManager stopped 16/04/29 10:24:29 INFO BlockManagerMaster: BlockManagerMaster stopped 16/04/29 10:24:30 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped! 16/04/29 10:24:30 INFO SparkContext: Successfully stopped SparkContext 16/04/29 10:24:30 INFO RemoteActorRefProvider$RemotingTerminator: Shutting down remote daemon. 16/04/29 10:24:30 INFO RemoteActorRefProvider$RemotingTerminator: Remote daemon shut down; proceeding with flushing remote transports. 16/04/29 10:24:30 INFO RemoteActorRefProvider$RemotingTerminator: Remoting shut down. 总结 java的代码实现spark API虽然代码冗余很多,但是很清楚显示了spark的执行过程,先比于scala的代码,较为清楚,而且java的代码和其他的项目结合效果可能好些。

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

分区索引的应用和实践 - 阿里云RDS PostgreSQL最佳实践

标签 PostgreSQL , partial index , partition index 背景 当表很大时,大家可能会想到分区表的概念,例如用户表,按用户ID哈希或者范围分区,拆成很多表。 又比如行为数据表,可以按时间分区,拆成很多表。 拆表的好处: 1、可以将表放到不同的表空间,表空间和块设备挂钩,例如历史数据访问量低,数据量大,可以放到机械盘所在的表空间。而活跃数据则可以放到SSD对应的表空间。 2、拆表后,方便维护,例如删除历史数据,直接DROP TABLE就可以了,不会产生REDO。 索引实际上也有分区的概念,例如按USER ID HASH分区,按时间分区等。 分区索引的好处与分区表的好处类似。同时还有其他好处: 1、不需要被检索的部分数据,可以不对它建立索引。 例如一张用户表,我们只检索已激活的用户,对于未激活的用户,我们不对它进行检索

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

【最佳实践实践总结 阿里云Elasticsearch 智能化运维思路

作者:梵寞—阿里云 Elasticsearch 团队高级开发工程师 阿里云 Elasticsearch 作为一个开箱即用的搜索引擎,其丰富的功能和极低的使用门槛吸引着越来越多的公司和用户选择它作为搜索和数据分析的工具。用户在运维 Elasticsearch 集群时往往会遇到很多难题,具体来说有下面列举的几点: 使用方式往往比较粗糙,默认的设置并不适合每一个集群和业务,非精细化的设计将会极大的增加集群隐患; 集群出现问题,无法及时定位原因、寻找解决方案,低效的沟通或者解决问题的方式可能会使得问题变得愈发严重; Elasticsearch 提供的监控指标繁杂,指标多,意义不明确,需要一定的专业知识才可以理解,缺乏全局视角; 此外,集群潜在的异常无法发现,更不能及时规避风险。 随着越来越多的用户选择使用阿里云 Elasticsearch 服务来支持

资源下载

更多资源
Mario

Mario

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

Nacos

Nacos

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

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

用户登录
用户注册