首页 文章 精选 留言 我的

精选列表

搜索[自动化研究],共10000篇文章
优秀的个人博客,低调大师

领导让我研究 Eureka 源码 | 启动过程

Eureka 源码之启动过程 大家好,我悟空。 最近在倒腾 Eureka 源码,因为大环境太卷了,必须得卷点源码才行,另外呢,能够读懂开源项目的源码、解决项目中遇到的问题是实力的象征,是吧?如果只是会用些中间件,那是不够的,和 CRUD 区别不大。 话不多说,源码走起。本篇是 Eureka 源码分析的开篇,后续会持续分享源码解析的文章。 首先呢,Eureka 服务的启动入口在这里:EurekaBootStrap.java 的 contextInitialized 方法。 关于源码的获取直接到官网下载就好了。https://github.com/Netflix/eureka 本文已收录到我的 github:https://github.com/Jackson0714/PassJava-Learning 一、初始化环境 启动类是 EurekaBootStrap.java,在这个路径下: \eureka\eureka-core\src\main\java\com\netflix\eureka\EurekaBootStrap.java 启动时序图: 启动代码: @Override public void contextInitialized(ServletContextEvent event) { initEurekaEnvironment(); initEurekaServerContext(); // 省略非核心代码 } 初始化环境的方法 initEurekaEnvironment(),点进去看下这个方法做了什么。 String dataCenter = ConfigurationManager.getConfigInstance().getString(EUREKA_DATACENTER); 就是获取配置管理类的一个单例。单例的实现方法用的是 双重检测 + volatile public static AbstractConfiguration getConfigInstance() { if (instance == null) { synchronized (ConfigurationManager.class) { if (instance == null) { instance = getConfigInstance(false)); } } } return instance; } instance 变量定义成了 volatile,保证可见性。 static volatile AbstractConfiguration instance = null; 线程 A 修改后,会将变量的值刷到主内存中,线程 B 会将主内存中的值刷回到自己的线程内存中,也就是说线程 A 改了后,线程 B 可以看到改了后的值。 可以参考之前我写的文章:反制面试官 - 14 张原理图 - 再也不怕被问 volatile 启动的整体流程如图: 二、初始化上下文 初始化上下文的时序图如下: 还是在 EurekaBootStrap.java 类中 contextInitialized 方法中,第二步调用了 initEurekaServerContext() 方法。 initEurekaServerContext 里面主要的操作分为六步: 第一步就是加载 配置文件。 2.1 加载 eureka-server 配置文件 基于接口的方式,获取配置项。 initEurekaServerContext 方法创建了一个 eurekaServerConfig 对象: EurekaServerConfig eurekaServerConfig = new DefaultEurekaServerConfig(); EurekaServerConfig 是一个接口,里面定义了很多获取配置项的方法。和定义常量来获取配置项的方式不同。比如获取 AccessId 和 SecretKey。 String getAWSAccessId(); String getAWSSecretKey(); 还有另外一种获取配置项的方式:Config.get(Constants.XX_XX),这种方式和上面的接口的方式相比: 常量的方式较容易取错变量。因为常量的定义都是大写,很可能拿到 XX_XY 变量,而接口的方法是驼峰命名的,更容易辨识,对于相似的变量,取一个辨识度更高的方法名即可。 常量的方式不易于修改。假如修改了常量名称,则需要全局搜索用到的地方,都改掉。如果是用接口的方式,则只需要修改接口方法中引用常量的地方即可,对于调用接口方法的地方是透明的。 2.1.1 创建默认的 eureka server 配置 new DefaultEurekaServerConfig (),会创建出一个默认的 server 配置,构造方法会调用 init 方法: public DefaultEurekaServerConfig() { init(); } 2.2.2 加载配置文件 private void init() { String env = ConfigurationManager.getConfigInstance().getString( EUREKA_ENVIRONMENT, TEST); ConfigurationManager.getConfigInstance().setProperty( ARCHAIUS_DEPLOYMENT_ENVIRONMENT, env); String eurekaPropsFile = EUREKA_PROPS_FILE.get(); try { // ConfigurationManager // .loadPropertiesFromResources(eurekaPropsFile); ConfigurationManager .loadCascadedPropertiesFromResources(eurekaPropsFile); } catch (IOException e) { logger.warn( "Cannot find the properties specified : {}. This may be okay if there are other environment " + "specific properties or the configuration is installed with a different mechanism.", eurekaPropsFile); } } 前两行是设置环境名称,后面几行是关键语句:获取配置文件,并放到 ConfigurationManager 单例中。 来看下 EUREKA_PROPS_FILE.get(); 做了什么。首先 EUREKA_PROPS_FILE 是这样定义的: private static final DynamicStringProperty EUREKA_PROPS_FILE = DynamicPropertyFactory .getInstance().getStringProperty("eureka.server.props", "eureka-server"); 用单例工厂 DynamicPropertyFactory 设置了默认值 eureka-server,然后 EUREKA_PROPS_FILE.get() 就会从缓存里面这个默认值。 然后再调用 loadCascadedPropertiesFromResources 方法,来加载配置文件。 首先会拼接默认的配置文件: String defaultConfigFileName = configName + ".properties"; 然后获取默认配置文件的配置项: Properties props = getPropertiesFromFile(url); 然后再拼接当前环境的配置文件 String envConfigFileName = configName + "-" + environment + ".properties"; 然后获取环境的配置文件的配置项并覆盖之前的默认配置项。 props.putAll(envProps); putAll 方法就是将这些属性放到一个 map 中。 然后这些配置项统一都交给 ConfigurationManager 来管理: config.loadProperties(props); 其实就是加载这个文件: 打开这个文件后,发现里面有几个 demo 配置项,不过都被注释了。 2.1.3 真正的配置项在哪? 上面可以看到 eureka-server.properties 都是空的,那配置项都配置在哪呢? 我们之前说过,DefaultEurekaServerConfig 是实现了 EurekaServerConfig 接口的,如下所示: public class DefaultEurekaServerConfig implements EurekaServerConfig 在 EurekaServerConfig 接口里面定义很多 get 方法,而 DefaultEurekaServerConfig 实现了这些 get 方法,来看下怎么实现的: @Override public int getWaitTimeInMsWhenSyncEmpty() { return configInstance.getIntProperty( namespace + "waitTimeInMsWhenSyncEmpty", (1000 * 60 * 5)).get(); } 里面的类似这样的 getXX 的方法,都有一个 default value,比如上面的是 1000 * 60 * 5,所以我们可以知道,配置项是在 DefaultEurekaServerConfig 类中定义的。 configInstance 这个单例又是 DynamicPropertyFactory 类型的,而在创建 configInstance 单例的时候,ConfigurationManager 还做了一些事情:将配置文件中的配置项放到 DynamicPropertyFactory 单例中,这样的话,DefaultEurekaServerConfig 中的 get 方法就可以获取到配置文件中的配置项了。具体的代码在 DynamicPropertyFactory 类中的 initWithConfigurationSource 方法中。 结合上面的加载配置文件的分析,可以得出结论:如果配置文件中没有配置,则用 DefaultEurekaServerConfig 定义的默认值。 2.1.4 加载配置文件小结 (1)创建一个 DefaultEurekaServerConfig 对象,实现了 EurekaServerConfig 接口,里面有很多获取配置项的方法。 (2)DefaultEurekaServerConfig 构造函数中调用了 init 方法。 (3)init 方法会加载 eureka-server.properties 配置文件,把里面的配置项都放到一个 map 中,然后交给 ConfigurationManager 来管理。 (4)DefaultEurekaServerConfig 对象里面有很多 get 方法,里面通过 hard code 定义了配置项的名称,当调用 get 方法时,调用的是 DynamicPropertyFactory 的获取配置项的方法,这些配置项如果在配置文件中有,则用配置项的。配置文件中的配置项是通过 ConfigurationManager 赋值给 DynamicPropertyFactory 的。 (5)当要获取配置项时,就调用对应的 get 方法,如果配置文件没有配置,则用默认值。 2.2 构造实例信息管理器 2.2.1 初始化服务实例的配置 instanceConfig 创建了一个 ApplicationInfoManager 对象,服务配置管理器,Application 可以理解为一个 Eureka client,作为一个应用程序向 Eureka 服务注册的。 applicationInfoManager = new ApplicationInfoManager( instanceConfig, new EurekaConfigBasedInstanceInfoProvider(instanceConfig).get()); 创建这个对象时,传了 instanceConfig,这个就是 eureka 实例的配置。这个 instanceConfig 和之前讲过的 EurekaServerConfig 很像,都是实现了一个接口,通过接口的 getXX 方法来获取配置信息。 2.2.2 构造服务实例 instanceInfo 另外一个参数是 EurekaConfigBasedInstanceInfoProvider,这个 Provider 是用来构造 instanceInfo(服务实例)。 怎么构造出来的呢?用到了设计模式中的构造器模式,而用到的配置信息就是从 EurekaInstanceConfig 里面获取到的。 InstanceInfo.Builder builder = InstanceInfo.Builder.newBuilder(vipAddressResolver); builder.setXX ... instanceInfo = builder.build(); setXX 的代码如下所示: 2.2.3 小结 (1)初始化服务实例的配置 instanceConfig 。 (2)用构造器模式初始化服务实例 instanceInfo。 (3)将 instanceConfig 和 instanceInfo 传给了 ApplicationInfoManager,交由它来管理。 2.3 初始化 eureka-client 2.3.1 初始化 eureka-client 配置 eurekaClient 是包含在 eureka-server 服务中的,用来跟其他 eureka-server 进行通信的。为什么还会有其他 eureka-server,因为在集群环境中,是会有多个 eureka 服务的,而服务之间是需要相互通信的。 初始化 eureka-client 代码: EurekaClientConfig eurekaClientConfig = new DefaultEurekaClientConfig(); eurekaClient = new DiscoveryClient(applicationInfoManager, eurekaClientConfig); 第一行又是初始化了一个配置,和之前初始化 server config,instance config 的地方很相似。也是通过接口方法里面的 DynamicPropertyFactory 来获取配置项的值。 eureka-client 也有一个加载配置文件的方法: Archaius1Utils.initConfig(CommonConstants.CONFIG_FILE_NAME); 这个文件就是 eureka-client.properties。 初始化配置的时候还初始化了一个 DefaultEurekaTransportConfig(),可以理解为传输的配置。 2.3.2 初始化 eurekaClient 再来看下第二行代码,创建了一个 DiscoveryClient 对象,赋值给了 eurekaClient。 创建 DiscoveryClient 对象的过程非常复杂,我们来细看下。 (1)拿到 eureka-client 的 config 、transport 的 config、instance 实例信息。 (2)判断是否要获取注册表信息,默认会获取。 if (config.shouldFetchRegistry()) 如果在配置文件中定义了 fetch-registry: false,则不会获取,单机 eureka 情况下,配置为 false,因为自己就包含了注册表信息,而且也不需要从其他 eureka 实例上获取配置信息。当在集群环境下,才需要获取注册表信息。 (3)判断是否要把自己注册到其他 eureka 上,默认会注册。 if (config.shouldRegisterWithEureka()) 单机情况下,配置 register-with-eureka: false。 (4)创建了一个支持任务调度的线程池。 scheduler = Executors.newScheduledThreadPool(2, new ThreadFactoryBuilder() .setNameFormat("DiscoveryClient-%d") .setDaemon(true) .build()); (5)创建了一个支持心跳检测的线程池。 heartbeatExecutor = new ThreadPoolExecutor( 1, clientConfig.getHeartbeatExecutorThreadPoolSize(), 0, TimeUnit.SECONDS, new SynchronousQueue<Runnable>(), new ThreadFactoryBuilder() .setNameFormat("DiscoveryClient-HeartbeatExecutor-%d") .setDaemon(true) .build() ); // use direct handoff (6)创建了一个支持缓存刷新的线程池。 cacheRefreshExecutor = new ThreadPoolExecutor( 1, clientConfig.getCacheRefreshExecutorThreadPoolSize(), 0, TimeUnit.SECONDS, new SynchronousQueue<Runnable>(), new ThreadFactoryBuilder() .setNameFormat("DiscoveryClient-CacheRefreshExecutor-%d") .setDaemon(true) .build() ); // use direct handoff (7)创建了一个支持 eureka client 和 eureka server 进行通信的对象 eurekaTransport = new EurekaTransport(); (8)初始化调度任务 initScheduledTasks(); 这个里面就会根据 fetch-registry 来判断是否需要开始调度执行刷新注册表信息,默认 30 s 调度一次。这个刷新的操作是由一个 CacheRefreshThread 线程来执行的。 同样的,也会根据 register-with-eureka 来判断是否需要开始调度执行发送心跳,默认 30 s 调度一次。这个发送心跳的操作由一个 HeartbeatThread 线程来执行的。 然后还创建了一个实例信息的副本,用来将自己本地的 instanceInfo 实例信息传给其他服务。什么时候发送这些信息呢? 又创建了一个监听器 statusChangeListener,这个监听器监听到状态改变时,就调用副本的 onDemandUpdate() 方法,将 instanceInfo 传给其他服务。 2.4 处理注册相关的流程 2.4.1 注册对象 创建了一个 PeerAwareInstanceRegistryImpl 对象,通过名字可以知道是可以感知集群实例注册表的实现类。通过官方注释可以知道这个类的作用: 处理所有的拷贝操作到其他节点,让他们保持同步。复制的操作包含 注册,续约,摘除,过期和状态变更。 当 eureka server 启动后,它尝试着从集群节点去获取所有的注册信息。如果获取失败了,当前 eureka server 在一段时间内不会让其他应用获取注册信息,默认 5 分钟。 自我保护机制:如果应用丢失续约的占比在一定时间内超过了设定的百分比,则 eureka 会报警,然后停止执行过期应用。 registry = new PeerAwareInstanceRegistryImpl( eurekaServerConfig, eurekaClient.getEurekaClientConfig(), serverCodecs, eurekaClient ); PeerAwareInstanceRegistryImpl 继承 AbstractInstanceRegistry 抽象类,构造函数主要做了以下事情: 初始化 server config 和 client config 的配置信息。 this.serverConfig = serverConfig; this.clientConfig = clientConfig; 初始化摘除的队列,队列长度为 1000。 this.recentCanceledQueue = new CircularQueue<Pair<Long, String>>(1000); 初始化注册的队列。 this.recentRegisteredQueue = new CircularQueue<Pair<Long, String>>(1000); 2.5 初始化上下文 2.5.1 集群节点帮助类 创建了一个 PeerEurekaNodes,它是一个帮助类,来管理集群节点的生命周期。 PeerEurekaNodes peerEurekaNodes = getPeerEurekaNodes( registry, eurekaServerConfig, eurekaClient.getEurekaClientConfig(), serverCodecs, applicationInfoManager ); 2.5.2 默认上下文 创建了一个 DefaultEurekaServerContext 默认上下文。 serverContext = new DefaultEurekaServerContext( eurekaServerConfig, serverCodecs, registry, peerEurekaNodes, applicationInfoManager ); 2.5.3 创建上下文的持有者 创建了一个 holder,用来持有上下文。其他地方想要获取上下文,就通过 holder 来获取。用到了单例模式。 EurekaServerContextHolder.initialize(serverContext); holder 的 initialize() 初始化方法是一个线程安全的方法。 public static synchronized void initialize(EurekaServerContext serverContext) { holder = new EurekaServerContextHolder(serverContext); } 定义了一个静态的私有的 holder 变量 private static EurekaServerContextHolder holder; 其他地方想获取 holder 的话,就通过 getInstance() 方法来获取 holder。 public static EurekaServerContextHolder getInstance() { return holder; } 然后想要获取上下文的就调用 holder 的 getServerContext() 方法。 public EurekaServerContext getServerContext() { return this.serverContext; } 2.5.4 初始化上下文 调用 serverContext 的 initialize() 方法来初始化。 public void initialize() throws Exception { logger.info("Initializing ..."); peerEurekaNodes.start(); registry.init(peerEurekaNodes); logger.info("Initialized"); } peerEurekaNodes.start(); 这个里面就是启动了一个定时任务,将集群节点的 URL 放到集合里面,这个集合不包含本地节点的 url。每隔一定时间,就更新 eureka server 集群的信息。 registry.init(peerEurekaNodes); 这个里面会初始化注册表,将集群中的 注册信息获取下,然后放到注册表里面。 2.6 其他 2.6.1 从相邻节点拷贝注册信息 int registryCount = registry.syncUp(); 2.6.2 eureka 监控 EurekaMonitors.registerAllStats(); 2.7 编译报错的解决方案 1、异常1 An exception occurred applying plugin request [id: 'nebula.netflixoss', version: '3.6.0'] 解决方案 plugins { id 'nebula.netflixoss' version '5.1.1' } 异常 2、 eureka-server-governator Plugin with id 'jetty' not found. 参考 https://blog.csdn.net/Sino_Crazy_Snail/article/details/79300058 三、总结 来一份 Eureka 启动的整体流程图 写了两本 PDF,回复 分布式 或 PDF 下载。 我的知识星球已开通,回复 知识星球 加入。 我是悟空,努力变强,变身超级赛亚人!

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

深入研究Apache Spark 3.0的新功能

直播回放:https://developer.aliyun.com/live/2894 以下是直播内容精华整理。 Spark3.0解决了超过3400个JIRAs,历时一年多,是整个社区集体智慧的成果。Spark SQL和Spark Cores是其中的核心模块,其余模块如PySpark等模块均是建立在两者之上。Spark3.0新增了太多的功能,无法一一列举,下图是其中24个相对来说比较重要的新功能,下文将会围绕这些进行简单介绍。 一、Performance 与性能相关的新功能主要有: Adaptive Query Execution Dynamic Partition Pruning Query Complication Speedup Join Hints (一)Adaptive Query Execution Adaptive Query Execution(AQE)在之前的版本里已经有所实现,但是之前的框架存在一些缺陷,导致使用不是很多,在Spark3.0中Databricks(Spark初创团队创建的大数据与AI智能公司)和Intel的工程师合作,解决了相关的问题。 在Spark1.0中所有的Catalyst Optimizer都是基于规则 (rule) 优化的。为了产生比较好的查询规则,优化器需要理解数据的特性,于是在Spark2.0中引入了基于代价的优化器 (cost-based optimizer),也就是所谓的CBO。然而,CBO也无法解决很多问题,比如: 数据统计信息普遍缺失,统计信息的收集代价较高; 储存计算分离的架构使得收集到的统计信息可能不再准确; Spark部署在某一单一的硬件架构上,cost很难被估计; Spark的UDF(User-defined Function)简单易用,种类繁多,但是对于CBO来说是个黑盒子,无法估计其cost。 总而言之,由于种种限制,Spark的优化器无法产生最好的Plan。也正是因为上诉原因,运行期的自适应调整就变得相当重要,对于Spark更是如此,于是有了AQE,其基本方法也非常简单易懂。如下图所示,在执行完部分的查询规划后,Spark可以收集到结果的统计信息,然后利用这些信息再对查询规划重新进行优化。这个优化的过程不是一次性的,而是自适应的,也就是说随着查询规划的执行会不断的进行优化, 而且尽可能地复用了现有优化器的已有优化规则。让整个查询优化变得更加灵活和自适应。 Spark3.0中AQE包括三个主要的运行期自适应功能: 可以基于运行期的统计信息,将Sort Merge Join 转换为Broadcast Hash Join; 可以基于数据运行中间结果的统计信息,减少reducer数量,避免数据在shuffle期间的过量分区导致性能损失; 可以处理数据分布不均导致的skew join。 更多的信息大家可以通过搜索引擎查询了解。 如果你是一个Spark的资深用户,可能你读了很多的调优宝典,其中第一条就是让你的Join变得更快的方法就是尽可能地使用Broadcast Hash Join。比如你可以增加spark.sql.autoBroadcastJoinThreshold 阈值,或者使用 broadcast HINT。但是这基本上属于艺高人胆大。首先,这种方法很难调,一不小心就会Out of Memory,甚至性能变得更差,即使现在产生了一定效果,但是随着负载的变化可能调优会完全失败。 也许你会想:Spark为什么不解决这个问题呢?这里有很多挑战,比如: 统计信息的缺失,统计信息的不准确,那么就是默认依据文件大小来预估表的大小,但是文件往往是压缩的,尤其是列存储格式,比如parquet 和 ORC,而Spark是基于行处理,如果数据连续重复,file size可能和真实的行存储的真实大小,差别非常之大。这也是为何提高autoBroadcastJoinThreshold,即使不是太大也可能会导致out of memory; Filter复杂、UDFs的使用都会使Spark无法准确估计Join输入数据量的大小。当你的query plan异常大和复杂的时候,这点尤其明显。 其中,Spark3.0中基于运行期的统计信息,将Sort Merge Join 转换为Broadcast Hash Join的过程如下图所示。 也许你还会看到调优宝典告诉你调整shuffle产生的partitions的数量。而当前默认数量是200,但是这个200为什么就不得而知了。然而,这个值设置为多少都不是最优的。其实在不同shuffle,数据的输入大小和分布绝大多数都是不一样。那么简单地用一个配置,让所有的shuffle来遵循,显然是不好的。要设得太小,每个partition的大小就会太大,那么GC的压力就会很大,aggregation和sort会更有可能的去spill数据到磁盘。但是,要是设太大,partition的大小就会太小,partition的数量会大。这个会导致不必要的IO,也让task调度器的压力剧增。那么调度器会导致所有task都变慢。这一系列问题在query plan复杂的时候变得尤为突出,还可能会影响到其他性能,最后耗时耗力却调优失败。 对于这个问题的解决,AQE就有优势了。如下图所示,AQE可以在运行期动态的调整partition来达到性能最优。 此外,数据分布不均是Spark调优的一个疑难杂症,它的表现有多种,比如若干task停滞不前,像是出现了bugs,又比如大量的disk spilling会导致很多节点都无事可做。此外,你也许会看到out of memory这种异常。其解决方法也很多,比如找到skew values然后重写query,或者在join的情况下增加skew keys来消除数据分布不均,但是无论哪种方法,都非常浪费时间,且后期难以维护。AQE解决问题的方式如下,其通过shuffle落地后的中间数据结果判断哪些partition是skew的,如果partition过大,就将其分成若干较小的partition,通过分而治之,总体性能大幅提升。 AQE的发布可以说是一个时代的开始,未来将会更进一步发展,引入更多自适应规则,让Spark可以随着数据分布和特性的变化自动改变Query plan,让更多的query编译静态优化变成运行时的动态优化。 (二)Dynamic Partition Pruning Dynamic Partition Pruning也是一个运行时的动态优化方法,简单来说就是我们可以通过Query的某些分支的中间结果来避免不必要的partition读取,这种方法是无法通过编译期推测出来的,只能在运行时根据结果来判断,这种方法对数据仓库的star-schema效果非常明显,在TPC-DS获得了非常明显的加速,可以加速2-18倍。 (三)Join Hints Join Hints是一个非常普遍的数据库的优化策略,在3.0之前已经有了Broadcast hash join,3.0之后的版本加了Sort-merge join、Shuffle hash join和 Shuffle nested loop join,但是要注意谨慎使用,因为数据的特性不同,很难保证一直有效,即使有效,也不代表一直有效,随着时间的变化,你的数据变了,可能会让你的query 变慢,变得不稳定。总体来说上面的四种Join的适用条件和特点如下所示,总而言之,使用Join Hints要谨慎。 二、Richer APIs Spark3.0简化了开发,不但增加了更多的新功能,也改善了众多现有的功能,让更多的用法成为可能,主要有: Accelerator-aware Scheduler Built-in Functions pandas UDF enhancements DELETE/UPDATE/MERGE in Catalyst (一)pandas UDF enhancements pandas UDF应该说是PySPark用户中最喜爱的特性之一,对于其功能和性能的提升应该都是喜闻乐见的,其发展历程如下图所示。 最新的pandas UDF和之前的不同之处在于引入了Python Type Hints,现在用户可以使用pandas中的数据类型比如pandas.Series等来表示pandas UDF的种类,不再需要记住原来的UDF类型,只需要指定正确的输入和输出类型即可。此外,pandas UDF可以分为pandas UDF和pandas API。 (二)Accelerator-aware Scheduler Accelerator-aware Scheduler是加速器的调度支持,狭义上也就是指GPU调度支持。加速器经常用来对特定负载做加速,目前,用户还是需要指定什么应用需要加速器资源,但是在将来我们会支持job或者stage级别的调度。Spark3.0中我们已经支持大多调度器,此外,我们还可以通过Web UI来监控GPU的使用,欢迎大家使用,更多详细资料大家可以到社区学习。 (三)Built-in Functions 为了让Spark3.0更方便实用,Spark社区按照其他的主流,比如数据库厂商等,内嵌了如上图所示的32个常用函数,这样用户就无须自己写UDF,并且速度更快。比如针对map类型,Spark3.0新增加了map_keys和map_values,更加地方便易用。其他新增加的更多内嵌函数大家可以到社区具体了解。 三、Monitoring and Debuggability Spark3.0也增加了一些对监控和调优的改进,主要有: Structured Streaming UI DDL/DML Enhancements Observable Metrics Event Log Rollover (一)Structured Streaming UI Structured Streaming是在Spark2.0中发布的,在Spark3.0中加入了UI的配置。新的UI主要包括了两种统计信息:已完成的Streaming查询聚合信息和未完成的Streaming查询的当前信息,包括Input Rate、Process Rate、Batch Duration和Operate Duration。 (二)DDL/DML Enhancements 我们还增加了各种DDL/DML命令,比如EXPLAIN和。EXPLAIN是性能调优的必备工具,读取EXPLAIN是每个用户的基本功,但是随着系统的运行,EXPLAIN的信息越来越多,而且信息多元、多样,在新的版本中我们引入了新的FORMATTED模式,如下所示,在开头处有一个非常精简的树状图,且之后的每个部分都有很详细的解释,更容易加更多的注意,这就从水平扩展变成了垂直扩展,更加的直观。 (三)Observable Metrics 我们还引入了Observable Metrics用以观测数据的质量。要知道数据质量对于很多Spark应用都是相当重要的,通常定义数据质量的Metrics还是非常容易的,比如用一些聚合参数,但是算出这个Metrics的值就非常麻烦,尤其对于流计算来说。 四、SQL Compatibility SQL兼容性也是Spark必不可提的话题,良好的兼容性更方便用户迁移到Spark平台,在Spark3.0中新增的主要功能有: ANSI Store Assignment Overflow Checking Reserved Keywords in Parser Proleptic Gregorian Calendar 也就是说,这个版本中我们让insert遵守了ANSI Store Assignment,并且增加了运行时的overflow的检查,还提供了一个模式让SQL Parser来准确地遵守ANSI标准的保留字,还切换了Calendar,这样更加符合ANSI的SQL标准。比如说我们想要插入两列数据,类型是int和string,如果将int插入到了string中,还是可以的,不会发生数据精度的损失和数据丢失;但是如果我们尝试将string类型插入到int类型中,就有可能发生数据损失甚至丢失。ANSI Store Assignment+Overflow Checking在输入不合法的时候就会在运行时抛出异常,需要注意的是这个设置默认是关闭的,可以根据个人需要打开。 五、Built-in Data Sources 在这个版本中我们提升了预装的数据源,比如Parquet table,我们可以对Nested Column做Column Pruning和Filter Pushdown,此外还支持了对CSV的Filter Pushdown,还引入了Binary Data Source来处理类似于二进制的图片文件。 六、Extensibility and Ecosystem Spark3.0继续加强了对生态圈的建设: 对Data Source V2 API的持续改善和catalog支持; 支持Java 11; 支持Hadoop 3; 支持Hive 3。 (一)Data Source V2 API+Catalog Support Spark3.0加上了对Catalog的支持来扩展Data Source API。Catalog plugin API可以让用户注册自己的catalog来实现对元数据的处理,这样可以让Spark用户更简单方便的使用数据源的表。对于没有实现Catalog plugin的数据源,用户需要先注册每个外部数据源的表才能访问,但是实现了Catalog plugin API之后我们只需要注册Catalog,然后就可以直接远程访问和操作catalog的表。对于数据源的开发者来说,什么时候支Data Source V2 API呢?下面是几点建议: 不过这里需要注意,Data Source V2还不是很稳定,开发者可能在未来还需要调整相关API的实现。大数据的发展相当迅速,Spark3.0为了能更方便的部署,我们升级了对各个组件和环境版本的支持,但是要注意以下事项。 关于生态圈,这里要提一下Koalas,它是一个纯的Python库,用Spark实现了绝大部分的pandas API,让pandas用户除了可以处理小数据,也可以处理大数据。Koalas对于pandas用户来说可以将pandas的代码扩展到大数据处理,使得学习PySpark变得更简单;对于现有的PySpark用户来说,多了更多的选择,可以用pandas API来解决生产力问题。过去一年多,Koalas的下载量是惊人的,在pip的下载量单日已经超过了37000,而且还在不断增长,5月的下载量也达到了85万。Koalas的代码其实不多,主要是API的实现,执行还是由Spark来做,所以Spark性能的提升对于Koalas用户来说是直接受益的。Koalas的发布周期想当频密,目前已经有33个发布,欢迎大家下载使用。 如何读和理解Spark UI对大多数新用户来说是一个很大的挑战,尤其对SQL用户来说,在Spark3.0中我们增加了自己的UI文档https://spark.apache.org/docs/latest/web-ui.html并且增加了SQL Reference ,https://spark.apache.org/docs/latest/sql-ref.html等,更详细的文档使得用户上手Spark的时候更加容易,欢迎大家去试一试Spark3.0,感受Spark的强大。 关键词:Spark3.0、SQL、PySpark、Koalas、pandas、UDF、AQE 阿里巴巴开源大数据技术团队成立Apache Spark中国技术社区,定期推送精彩案例,技术专家直播,问答区近万人Spark技术同学在线提问答疑,只为营造纯粹的Spark氛围,欢迎钉钉扫码加入!

资源下载

更多资源
腾讯云软件源

腾讯云软件源

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

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应用均可从中受益。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

用户登录
用户注册