本文还有配套的精品资源,点击获取 menu-r.4af5f7ec.gif

简介:源码分析是提升IT技术水平、深入理解软件运行机制的重要方法。本资源为名为“xmind-master”的压缩包,包含使用Xmind工具整理的多个主流Java及大数据框架的源码分析思维导图,涵盖Spring、SpringMVC、MyBatis、Spark、Flink及JVM等核心技术。通过这些思维导图,开发者可系统理解各框架的架构设计、核心流程、组件交互及设计模式,掌握并发控制、异常处理与JVM优化技巧,全面提升代码质量与问题解决能力。
xmind:我平时做的写原始代码分析,包括spring,springmvc,mybatis,spark,flink等各种框架原始码以xmind的形式记录

1. Java主流框架源码分析概述

在Java生态中,Spring、SpringMVC、MyBatis、Spark、Flink等框架已成为企业级应用和大数据处理的核心支撑。理解这些框架的设计理念与源码结构,不仅有助于提升系统性能调优能力,也为深入掌握框架扩展机制、实现定制化开发打下坚实基础。

本章将从宏观视角出发,概述这些主流框架的基本架构与核心运行机制。例如,Spring框架通过IoC容器实现对象管理解耦,SpringMVC通过DispatcherServlet实现请求的统一调度,MyBatis则通过动态代理与SQL映射实现了灵活的ORM操作,而Spark和Flink分别以批处理和流式处理为核心构建了高性能的计算引擎。

通过源码分析,我们可以:

  • 深入理解框架内部组件协作机制;
  • 掌握关键设计模式和架构思想;
  • 提升排查复杂问题的能力;
  • 为框架二次开发与优化提供技术支撑。

接下来的章节将围绕这些框架的核心模块展开源码级深度剖析,逐步深入其实现细节。

2. Spring框架核心模块源码分析

Spring 框架作为 Java 生态中最为广泛使用的轻量级容器,其核心模块的设计与实现直接影响了整个框架的灵活性、可扩展性和性能。本章将深入分析 Spring 框架的四大核心模块:容器启动流程、Bean 生命周期管理、AOP 实现机制以及事务管理机制。通过源码分析,我们将逐步揭开 Spring 框架内部的运行原理,帮助开发者理解其设计思想与实现细节。

2.1 Spring容器的启动流程

Spring 容器的启动是整个应用运行的基础,其核心在于 ApplicationContext 的初始化过程。Spring 容器通过 ApplicationContext 接口来提供对 Bean 的管理能力,包括加载配置、实例化 Bean、依赖注入等。

2.1.1 ApplicationContext的初始化过程

Spring 容器的启动通常以 ClassPathXmlApplicationContext AnnotationConfigApplicationContext 为例进行分析。以 AnnotationConfigApplicationContext 为例,其构造函数会触发容器的初始化流程。

public class ApplicationContextExample {
    public static void main(String[] args) {
        ApplicationContext context = new AnnotationConfigApplicationContext(AppConfig.class);
    }
}

该构造函数的执行流程如下:

  1. 注册配置类 :将 AppConfig.class 注册为配置类,用于后续扫描 Bean。
  2. 创建 BeanFactory :初始化 DefaultListableBeanFactory ,这是 Spring 的核心容器。
  3. 加载 Bean 定义 :通过 AnnotatedBeanDefinitionReader 注册配置类中的 Bean 定义。
  4. 刷新容器 :调用 refresh() 方法,完成容器的初始化和 Bean 的加载。

refresh() 方法是整个容器启动的核心方法,其源码如下(简化版):

public void refresh() throws BeansException, IllegalStateException {
    prepareRefresh();
    ConfigurableListableBeanFactory beanFactory = obtainFreshBeanFactory();
    prepareBeanFactory(beanFactory);
    postProcessBeanFactory(beanFactory);
    invokeBeanFactoryPostProcessors(beanFactory);
    registerBeanPostProcessors(beanFactory);
    initMessageSource();
    initApplicationEventMulticaster();
    onRefresh();
    registerListeners();
    finishBeanFactoryInitialization(beanFactory);
    finishRefresh();
}
源码逻辑分析:
  • prepareRefresh() :准备刷新上下文,包括设置启动时间、激活状态等。
  • obtainFreshBeanFactory() :获取或创建一个新的 BeanFactory。
  • prepareBeanFactory(beanFactory) :配置 BeanFactory 的基本属性,如类加载器、表达式解析器等。
  • postProcessBeanFactory(beanFactory) :允许子类在 BeanFactory 配置完成后进行自定义修改。
  • invokeBeanFactoryPostProcessors(beanFactory) :调用所有 BeanFactoryPostProcessor ,用于修改 Bean 定义。
  • registerBeanPostProcessors(beanFactory) :注册所有 BeanPostProcessor ,用于在 Bean 创建过程中进行干预。
  • initMessageSource() :初始化消息源,用于国际化支持。
  • initApplicationEventMulticaster() :初始化事件广播器。
  • onRefresh() :模板方法,由子类实现,用于特定上下文刷新。
  • registerListeners() :注册所有 ApplicationListener。
  • finishBeanFactoryInitialization(beanFactory) :完成所有非延迟加载的单例 Bean 的初始化。
  • finishRefresh() :发布上下文刷新完成事件。

2.1.2 BeanFactory与FactoryBean的实现机制

BeanFactory 是 Spring 容器的核心接口,提供了获取 Bean、判断 Bean 是否存在、获取 Bean 类型等基本功能。

BeanFactory 接口核心方法:
public interface BeanFactory {
    Object getBean(String name) throws BeansException;
    <T> T getBean(String name, Class<T> requiredType) throws BeansException;
    boolean containsBean(String name);
    boolean isSingleton(String name) throws BeansException;
    boolean isPrototype(String name) throws BeansException;
    String[] getAliases(String name);
}

BeanFactory 是 Spring 的最基础容器,而 FactoryBean 是用于创建复杂对象的工厂类。它与普通 Bean 的区别在于, FactoryBean 本身是一个 Bean,但其返回的是 getObject() 方法的返回值。

FactoryBean 示例:
public class MyFactoryBean implements FactoryBean<MyObject> {
    @Override
    public MyObject getObject() throws Exception {
        return new MyObject("Created by FactoryBean");
    }

    @Override
    public Class<?> getObjectType() {
        return MyObject.class;
    }

    @Override
    public boolean isSingleton() {
        return true;
    }
}

在 Spring 配置中注册该 FactoryBean:

<bean id="myFactoryBean" class="com.example.MyFactoryBean"/>

当调用 context.getBean("myFactoryBean") 时,返回的是 MyObject 实例,而非 MyFactoryBean 本身。

源码分析:

AbstractBeanFactory 中, getBean() 方法会判断是否为 FactoryBean ,并调用其 getObject() 方法生成实际 Bean。

protected Object getObjectForBeanInstance(
        Object beanInstance, String name, String beanName, @Nullable RootBeanDefinition mbd) {

    if (isFactoryDereference(name) && !(beanInstance instanceof FactoryBean)) {
        throw new BeanIsNotAFactoryException(transformedBeanName(name), beanInstance.getClass());
    }

    if (!(beanInstance instanceof FactoryBean)) {
        return beanInstance;
    }

    FactoryBean<?> factory = (FactoryBean<?>) beanInstance;
    Object object = null;

    if (mbd != null) {
        mbd.isFactoryBean = true;
    }
    object = getObjectFromFactoryBean(factory, beanName, !synthetic);
    return object;
}
参数说明:
  • beanInstance :当前 Bean 的实例。
  • name :请求的 Bean 名称。
  • beanName :Bean 的实际名称。
  • mbd :Bean 的定义信息。

2.1.3 容器启动过程中的事件机制

Spring 提供了基于观察者模式的事件机制,用于在容器生命周期中发布和监听事件。

核心类与接口:
  • ApplicationEvent :所有事件的基类。
  • ApplicationListener :事件监听器接口。
  • ApplicationEventPublisher :事件发布接口。
  • ApplicationEventMulticaster :事件广播器,负责事件的注册与发布。
示例:自定义事件与监听器
public class CustomEvent extends ApplicationEvent {
    public CustomEvent(Object source) {
        super(source);
    }
}

public class CustomEventListener implements ApplicationListener<CustomEvent> {
    @Override
    public void onApplicationEvent(CustomEvent event) {
        System.out.println("Received custom event: " + event);
    }
}

在容器启动后发布事件:

((AbstractApplicationContext) context).publishEvent(new CustomEvent(context));
源码流程图:
graph TD
    A[ApplicationContext启动] --> B[注册监听器]
    B --> C[发布事件]
    C --> D[事件广播器收集监听器]
    D --> E[调用监听器onApplicationEvent方法]
表格:Spring 事件机制组件对比
组件 功能说明
ApplicationEvent 事件基类,用户可继承定义事件
ApplicationListener 监听器接口,用于处理事件
ApplicationEventPublisher 事件发布接口
ApplicationEventMulticaster 事件广播器,管理监听器并发布事件
源码分析:

AbstractApplicationContext 中, initApplicationEventMulticaster() 初始化事件广播器:

protected void initApplicationEventMulticaster() {
    ConfigurableListableBeanFactory beanFactory = getBeanFactory();
    if (beanFactory.containsLocalBean(APPLICATION_EVENT_MULTICASTER_BEAN_NAME)) {
        this.applicationEventMulticaster =
            beanFactory.getBean(APPLICATION_EVENT_MULTICASTER_BEAN_NAME, ApplicationEventMulticaster.class);
    } else {
        this.applicationEventMulticaster = new SimpleApplicationEventMulticaster(beanFactory);
        beanFactory.registerSingleton(APPLICATION_EVENT_MULTICASTER_BEAN_NAME, this.applicationEventMulticaster);
    }
}

通过 SimpleApplicationEventMulticaster ,Spring 可以将事件广播给所有注册的监听器。

(本章内容将持续扩展,后续将深入分析 Bean 生命周期、AOP 和事务机制的源码实现)

3. SpringMVC MVC架构源码流程解析

SpringMVC 是 Spring 框架中用于构建 Web 应用的核心模块,其基于 MVC(Model-View-Controller)架构,将业务逻辑、数据模型和用户界面分离,便于开发与维护。本章将深入剖析 SpringMVC 的核心处理流程,包括请求的接收、Controller 方法的调用、视图解析、响应构建,以及拦截器和异步支持等机制。通过源码层面的分析,帮助读者理解 SpringMVC 的内部工作原理,为性能优化、问题排查和框架定制打下基础。

3.1 DispatcherServlet 的请求处理流程

SpringMVC 的核心是 DispatcherServlet ,它继承自 HttpServlet ,负责接收所有 HTTP 请求并进行统一调度。整个处理流程由多个组件协同完成,包括 HandlerMapping ViewResolver HandlerAdapter 等。

3.1.1 请求的接收与处理线程模型

DispatcherServlet 的请求处理入口是 doService 方法,最终调用 doDispatch 方法进行请求分发。SpringMVC 使用多线程处理请求,每个请求由一个线程处理,保证并发性。

protected void doDispatch(HttpServletRequest request, HttpServletResponse response) throws Exception {
    HttpServletRequest processedRequest = checkMultipart(request);
    HandlerExecutionChain mappedHandler = null;
    ModelAndView mv = null;

    try {
        mappedHandler = getHandler(processedRequest); // 获取匹配的Handler
        HandlerAdapter ha = getHandlerAdapter(mappedHandler.getHandler()); // 获取适配器
        mv = ha.handle(processedRequest, response, mappedHandler.getHandler()); // 执行Controller方法
    } finally {
        processDispatchResult(processedRequest, response, mappedHandler, mv, dispatchException); // 处理结果
    }
}
逻辑分析
  • checkMultipart() :检查请求是否为文件上传,若是则进行封装处理。
  • getHandler() :遍历所有 HandlerMapping ,找到匹配当前请求的 Controller。
  • getHandlerAdapter() :根据 Handler 类型获取适配器,如 RequestMappingHandlerAdapter
  • ha.handle() :调用 Controller 方法,返回 ModelAndView
  • processDispatchResult() :处理结果,包括视图渲染和异常处理。
线程模型

SpringMVC 使用标准的 Servlet 线程模型,每个请求由一个独立线程处理。Tomcat 等 Web 容器默认使用线程池管理请求线程,提高并发处理能力。

3.1.2 HandlerMapping 的匹配机制

HandlerMapping 负责将请求 URL 映射到具体的 Controller 方法。SpringMVC 提供了多种实现,如 BeanNameUrlHandlerMapping RequestMappingHandlerMapping

典型流程
  1. 注册 Handler :在应用启动时, RequestMappingHandlerMapping 会扫描所有 @Controller 注解的类,并提取 @RequestMapping 的 URL 路径。
  2. 请求匹配 :在请求到达时, HandlerMapping 会根据 URL 匹配相应的 Controller 方法。
源码片段
public class RequestMappingHandlerMapping extends AbstractHandlerMethodMapping<RequestMappingInfo> {
    @Override
    protected Set<RequestMappingInfo> getMatchingMapping(RequestMappingInfo mapping, HttpServletRequest request) {
        return mapping.getMatchingCondition(request);
    }
}
说明
  • AbstractHandlerMethodMapping 是所有 HandlerMapping 的基类,提供统一的注册和匹配接口。
  • RequestMappingInfo 封装了 URL、HTTP 方法、参数等信息。
  • getMatchingCondition() 用于判断当前请求是否满足映射条件。
表格:常见 HandlerMapping 实现类
HandlerMapping 实现类 说明
BeanNameUrlHandlerMapping 根据 Bean 名匹配 URL
SimpleUrlHandlerMapping 配置方式指定 URL 与 Handler 的映射
RequestMappingHandlerMapping 支持 @RequestMapping 注解的 Controller 映射

3.1.3 Controller 方法的调用过程

Controller 方法的调用由 HandlerAdapter 完成,其核心是 handle() 方法。SpringMVC 支持多种类型的 Controller,如注解驱动的 @RequestMapping 方法和传统的 Controller 接口。

调用流程
  1. 参数解析 :通过 HandlerMethodArgumentResolver 解析方法参数,如 @RequestParam @RequestBody
  2. 方法调用 :使用反射机制调用 Controller 方法。
  3. 结果处理 :将返回值封装为 ModelAndView 或直接写入响应。
源码片段
public class RequestMappingHandlerAdapter extends AbstractHandlerMethodAdapter {
    @Override
    protected ModelAndView handleInternal(HttpServletRequest request, HttpServletResponse response, HandlerMethod handlerMethod) throws Exception {
        Object[] args = resolveHandlerArguments(handlerMethod, request, response);
        Object returnValue = handlerMethod.invokeForRequest(request, response, args);
        return getModelAndView(returnValue, handlerMethod);
    }

    private Object[] resolveHandlerArguments(HandlerMethod handlerMethod, HttpServletRequest request, HttpServletResponse response) {
        List<HandlerMethodArgumentResolver> resolvers = this.argumentResolvers;
        List<Object> args = new ArrayList<>();
        for (MethodParameter parameter : handlerMethod.getMethodParameters()) {
            for (HandlerMethodArgumentResolver resolver : resolvers) {
                if (resolver.supportsParameter(parameter)) {
                    args.add(resolver.resolveArgument(parameter, request, response));
                }
            }
        }
        return args.toArray();
    }
}
说明
  • resolveHandlerArguments() :遍历所有参数解析器,找到匹配的解析器解析参数。
  • invokeForRequest() :通过反射调用 Controller 方法。
  • getModelAndView() :根据返回值类型生成 ModelAndView
典型参数解析器
解析器类名 支持的参数类型
RequestParamMethodArgumentResolver @RequestParam
PathVariableMethodArgumentResolver @PathVariable
RequestBodyMethodArgumentResolver @RequestBody
ModelAttributeMethodArgumentResolver @ModelAttribute

3.2 视图解析与响应构建

SpringMVC 通过 ViewResolver 将逻辑视图名解析为实际的视图对象(如 JSP、Thymeleaf 模板),并使用 View 渲染模型数据,最终返回 HTTP 响应。

3.2.1 ModelAndView 的构建与处理

ModelAndView 是 Controller 返回值的封装,包含模型数据(Model)和视图名(View Name)。

典型流程
  1. Controller 返回 ModelAndView :如 return new ModelAndView("home", model);
  2. 解析视图 ViewResolver 根据视图名查找实际视图对象。
  3. 渲染模型 :调用 render() 方法将模型数据渲染到视图中。
  4. 写入响应 :视图渲染结果写入 HTTP 响应流。
源码片段
public interface ViewResolver {
    View resolveViewName(String viewName, Locale locale) throws Exception;
}

public class InternalResourceViewResolver implements ViewResolver {
    @Override
    public View resolveViewName(String viewName, Locale locale) throws Exception {
        return new InternalResourceView("/WEB-INF/views/" + viewName + ".jsp");
    }
}
说明
  • ViewResolver 接口定义了视图解析方法。
  • InternalResourceViewResolver 是常用的实现,用于 JSP 视图解析。

3.2.2 ViewResolver 的配置与实现

SpringMVC 支持多种视图解析器,开发者可根据需求配置。

配置方式(XML)
<bean id="viewResolver" class="org.springframework.web.servlet.view.InternalResourceViewResolver">
    <property name="prefix" value="/WEB-INF/views/" />
    <property name="suffix" value=".jsp" />
</bean>
典型实现类
ViewResolver 实现类 说明
InternalResourceViewResolver JSP 视图解析器
ThymeleafViewResolver Thymeleaf 模板引擎解析器
FreeMarkerViewResolver FreeMarker 模板解析器

3.2.3 响应数据的序列化与返回

对于 RESTful 接口,SpringMVC 支持直接返回 JSON 或 XML 数据。通常使用 @ResponseBody 注解或 ResponseEntity 对象。

示例代码
@GetMapping("/users")
@ResponseBody
public List<User> getAllUsers() {
    return userService.findAll();
}
底层处理流程
  1. 检测 @ResponseBody :若方法有该注解,则跳过视图解析。
  2. 使用 HttpMessageConverter :根据返回值类型选择合适的转换器,如 Jackson2JsonMessageConverter
  3. 写入响应体 :将对象序列化为 JSON 或 XML,并写入响应流。
常见 HttpMessageConverter
转换器类名 支持的数据类型
Jackson2JsonMessageConverter JSON
Jaxb2RootElementHttpMessageConverter XML
StringHttpMessageConverter String
ByteArrayHttpMessageConverter byte[]

3.3 SpringMVC 的拦截器机制

SpringMVC 提供了拦截器(Interceptor)机制,用于在请求处理的不同阶段执行公共逻辑,如权限验证、日志记录等。

3.3.1 拦截器的注册与执行顺序

拦截器通过 WebMvcConfigurer 接口进行注册,支持按 URL 模式匹配。

示例配置
@Configuration
@EnableWebMvc
public class WebConfig implements WebMvcConfigurer {
    @Override
    public void addInterceptors(InterceptorRegistry registry) {
        registry.addInterceptor(new AuthInterceptor())
                .addPathPatterns("/**")
                .excludePathPatterns("/login", "/error");
    }
}
拦截器接口
public interface HandlerInterceptor {
    default boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) throws Exception {
        return true;
    }

    default void postHandle(HttpServletRequest request, HttpServletResponse response, Object handler, ModelAndView modelAndView) throws Exception {
    }

    default void afterCompletion(HttpServletRequest request, HttpServletResponse response, Object handler, Exception ex) throws Exception {
    }
}
方法说明
方法名 说明
preHandle() 在 Controller 方法执行前调用,返回 false 表示中断请求
postHandle() 在 Controller 方法执行后调用,可修改 ModelAndView
afterCompletion() 在整个请求完成(视图渲染完成后)调用,可用于资源清理
拦截器执行流程图
graph TD
    A[DispatcherServlet] --> B[preHandle]
    B --> C[Controller Method]
    C --> D[postHandle]
    D --> E[ViewResolver]
    E --> F[View Render]
    F --> G[afterCompletion]

3.4 异步请求处理与 WebSocket 支持

SpringMVC 支持异步请求处理和 WebSocket 协议,适用于高并发场景和实时通信需求。

3.4.1 异步 Controller 的实现机制

SpringMVC 支持返回 Callable DeferredResult ,实现异步响应。

示例代码
@GetMapping("/async")
public Callable<String> asyncCall() {
    return () -> "Hello from async!";
}
异步处理流程
  1. 主线程提交任务 :Controller 返回 Callable ,由任务线程池执行。
  2. 任务完成回调 :任务完成后,由 SpringMVC 将结果写入响应。
  3. 释放主线程 :主线程释放,提高并发处理能力。
优势
  • 减少线程阻塞,提高系统吞吐量。
  • 适用于耗时操作,如远程调用、批量处理等。

3.4.2 WebSocket 的握手与消息处理流程

SpringMVC 支持 WebSocket 通信,提供 WebSocketHandler 接口和 @ServerEndpoint 注解。

示例代码
@Component
public class ChatWebSocketHandler implements TextWebSocketHandler {
    @Override
    public void handleTextMessage(WebSocketSession session, TextMessage message) {
        session.sendMessage(new TextMessage("Echo: " + message.getPayload()));
    }
}
握手流程
  1. 客户端发起 WebSocket 握手请求
  2. 服务端 WebSocketHttpRequestHandler 处理握手
  3. 建立 WebSocket 会话连接
  4. 双方通过 WebSocketSession 发送/接收消息
WebSocket 消息处理流程图
graph TD
    A[Client] --> B[WebSocket Handshake]
    B --> C[WebSocketHttpRequestHandler]
    C --> D[Session Created]
    D --> E[Message Handling]
    E --> F[TextWebSocketHandler]
    F --> G[Send/Receive Messages]

本章深入剖析了 SpringMVC 的核心处理流程,从请求的接收、Controller 方法的调用、视图解析与响应构建,到拦截器和异步处理机制,帮助读者全面理解 SpringMVC 的底层实现原理。下一章将进入 MyBatis 框架的源码分析,探讨其 ORM 实现与 SQL 映射机制。

4. MyBatis ORM实现与SQL映射源码分析

4.1 MyBatis核心组件的初始化流程

4.1.1 Configuration的加载与解析

MyBatis 的核心初始化流程从 Configuration 对象的创建开始。 Configuration 是 MyBatis 中非常重要的配置类,它包含了 MyBatis 的全局配置信息,如数据库连接池、事务管理器、类型别名、映射器配置等。

通常,MyBatis 的配置是通过 mybatis-config.xml 文件进行定义的。在初始化阶段,MyBatis 会通过 XMLConfigBuilder 来解析该配置文件,并将解析后的信息封装到 Configuration 对象中。

String resource = "mybatis-config.xml";
InputStream inputStream = Resources.getResourceAsStream(resource);
SqlSessionFactory sqlSessionFactory = new SqlSessionFactoryBuilder().build(inputStream);

在上述代码中, SqlSessionFactoryBuilder build 方法会调用 XMLConfigBuilder.parse() 方法,最终返回一个包含完整配置信息的 SqlSessionFactory

XMLConfigBuilder 的解析流程

XMLConfigBuilder 是 MyBatis 中用于解析主配置文件的核心类。其主要职责是解析 mybatis-config.xml 文件,并将其转换为 Configuration 实例。

解析流程如下:

  1. 读取 XML 文件 :通过 Document 对象读取 XML 配置文件。
  2. 解析 <configuration> 根节点 :依次解析 <properties> <settings> <typeAliases> <plugins> <environments> 等子节点。
  3. 构建 Configuration 对象 :将解析得到的配置信息填充到 Configuration 对象中。

例如,解析 <typeAliases> 节点时,MyBatis 会注册类型别名,以便在映射文件中使用更简洁的类名:

<typeAliases>
    <typeAlias alias="User" type="com.example.model.User"/>
</typeAliases>

对应的源码解析逻辑如下:

private void parseConfiguration(XNode root) {
    try {
        // 解析 properties 配置
        propertiesElement(root.evalNode("properties"));
        // 解析 settings 配置
        Properties settings = settingsAsProperties(root.evalNode("settings"));
        // 注册类型别名
        typeAliasesElement(root.evalNode("typeAliases"));
        // 解析插件
        pluginElement(root.evalNode("plugins"));
        // 解析环境配置
        environmentsElement(root.evalNode("environments"));
        // 解析映射器
        mapperElement(root.evalNode("mappers"));
    } catch (Exception e) {
        throw new BuilderException("Error parsing SQL Mapper Configuration. Cause: " + e, e);
    }
}

逻辑分析
- XNode 是 MyBatis 对 XML 节点的封装,用于简化 XML 解析。
- 每个配置节点都会被解析并转换为对应的 Java 对象,最终整合到 Configuration 中。
- 例如, mapperElement 方法会遍历 <mappers> 节点下的所有映射文件,并调用 XMLMapperBuilder 解析每个映射文件。

示例: <environments> 配置的解析

MyBatis 支持多环境配置,开发者可以在 <environments> 节点中定义多个数据库环境,并通过 default 属性指定默认环境。

<environments default="development">
    <environment id="development">
        <transactionManager type="JDBC"/>
        <dataSource type="POOLED">
            <property name="driver" value="com.mysql.cj.jdbc.Driver"/>
            <property name="url" value="jdbc:mysql://localhost:3306/mydb"/>
            <property name="username" value="root"/>
            <property name="password" value="123456"/>
        </dataSource>
    </environment>
</environments>

对应的解析逻辑如下:

private void environmentsElement(XNode context) throws Exception {
    if (context != null) {
        if (environment == null) {
            environment = context.getStringAttribute("default");
        }
        for (XNode child : context.evalNodes("environment")) {
            String id = child.getStringAttribute("id");
            Environment.Builder environmentBuilder = new Environment.Builder(id);
            // 解析事务管理器
            environmentBuilder.transactionFactory(transactionManagerElement(child.evalNode("transactionManager")));
            // 解析数据源
            environmentBuilder.dataSource(dataSourceElement(child.evalNode("dataSource")));
            configuration.addEnvironment(environmentBuilder.build());
        }
    }
}

参数说明
- id :环境标识符,用于区分不同环境。
- transactionManager :事务管理器类型,如 JDBC MANAGED
- dataSource :数据源类型,如 UNPOOLED POOLED JNDI

4.1.2 SqlSessionFactory的构建过程

SqlSessionFactory 是 MyBatis 的核心工厂类,用于创建 SqlSession 实例。它的构建依赖于前面解析得到的 Configuration 对象。

构建流程如下:

  1. 构建 SqlSessionFactory :通过 SqlSessionFactoryBuilder build 方法,将 Configuration 对象封装为 DefaultSqlSessionFactory
  2. 缓存映射器接口 :在 SqlSessionFactory 初始化过程中,会将所有的 Mapper 接口注册为动态代理类。
  3. 返回 SqlSessionFactory 实例 :最终返回的 SqlSessionFactory 实例可用于创建 SqlSession
public SqlSessionFactory build(Configuration config) {
    return new DefaultSqlSessionFactory(config);
}
DefaultSqlSessionFactory 的结构
public class DefaultSqlSessionFactory implements SqlSessionFactory {
    private final Configuration configuration;

    public DefaultSqlSessionFactory(Configuration configuration) {
        this.configuration = configuration;
    }

    @Override
    public SqlSession openSession() {
        return openSessionFromDataSource(configuration.getDefaultExecutorType(), null, false);
    }

    private SqlSession openSessionFromDataSource(ExecutorType execType, TransactionIsolationLevel level, boolean autoCommit) {
        Transaction tx = null;
        try {
            final Environment environment = configuration.getEnvironment();
            final TransactionFactory transactionFactory = getTransactionFactoryFromEnvironment(environment);
            tx = transactionFactory.newTransaction(environment.getDataSource(), level, autoCommit);
            final Executor executor = configuration.newExecutor(tx, execType);
            return new DefaultSqlSession(configuration, executor, autoCommit);
        } catch (Exception e) {
            closeTransaction(tx);
            throw ExceptionFactory.wrapException("Error opening session.  Cause: " + e, e);
        }
    }
}

逻辑分析
- openSession() 方法最终会调用 openSessionFromDataSource() 方法,用于创建一个数据库会话。
- Executor 是 MyBatis 的执行器,负责执行 SQL 语句、缓存处理等。
- DefaultSqlSession 是对数据库操作的封装,提供了执行 SQL、获取结果等功能。

4.1.3 Mapper接口的动态代理实现

MyBatis 使用动态代理机制将 Mapper 接口转换为可执行的 SQL 操作类。开发者只需要定义接口和注解(或 XML 映射文件),MyBatis 就会自动为其生成代理类。

动态代理流程
  1. 注册 Mapper 接口 :在 Configuration 初始化过程中,所有 Mapper 接口会被注册到 MapperRegistry 中。
  2. 生成代理类 :当调用 getMapper() 方法时,MyBatis 会使用 JDK 动态代理 CGLIB 生成代理类。
  3. 执行 SQL 语句 :代理类在调用方法时,会根据方法名和参数构建 SQL 语句,并调用底层 Executor 执行。
示例:Mapper 接口定义
public interface UserMapper {
    @Select("SELECT * FROM users WHERE id = #{id}")
    User selectById(int id);
}
获取 Mapper 实例
SqlSession sqlSession = sqlSessionFactory.openSession();
UserMapper userMapper = sqlSession.getMapper(UserMapper.class);
User user = userMapper.selectById(1);
源码解析: getMapper() 方法
public <T> T getMapper(Class<T> type) {
    return configuration.getMapper(type, this);
}
public <T> T getMapper(Class<T> type, SqlSession sqlSession) {
    return mapperRegistry.getMapper(type, sqlSession);
}
public <T> T getMapper(Class<T> type, SqlSession sqlSession) {
    final MapperProxyFactory<T> mapperProxyFactory = (MapperProxyFactory<T>) knownMappers.get(type);
    return mapperProxyFactory.newInstance(sqlSession);
}
MapperProxyFactory 的作用
public class MapperProxyFactory<T> {

    private final Class<T> mapperInterface;

    public MapperProxyFactory(Class<T> mapperInterface) {
        this.mapperInterface = mapperInterface;
    }

    protected T newInstance(MapperProxy<T> mapperProxy) {
        return (T) Proxy.newProxyInstance(mapperInterface.getClassLoader(), new Class[] { mapperInterface }, mapperProxy);
    }

    public T newInstance(SqlSession sqlSession) {
        return newInstance(new MapperProxy<>(sqlSession, mapperInterface, methodCache));
    }
}

逻辑分析
- MapperProxyFactory 使用 JDK 动态代理 创建代理类。
- MapperProxy 是 InvocationHandler 的实现类,负责拦截方法调用并生成 SQL。
- methodCache 缓存了方法与 SQL 映射的关系,提高性能。

MapperProxy invoke() 方法
public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
    try {
        if (Object.class.equals(method.getDeclaringClass())) {
            return method.invoke(this, args);
        } else {
            return cachedInvoker(method).invoke(proxy, method, args, sqlSession);
        }
    } catch (Throwable t) {
        throw ExceptionUtil.unwrapThrowable(t);
    }
}

逻辑分析
- 如果是 Object 类的方法(如 toString() hashCode() ),直接调用原方法。
- 否则,调用 cachedInvoker 获取 SQL 执行器并执行。

小结

MyBatis 的初始化流程围绕 Configuration SqlSessionFactory Mapper 动态代理三个核心组件展开。通过 XML 解析和动态代理机制,MyBatis 实现了灵活的 ORM 框架架构,为后续 SQL 执行和事务管理奠定了基础。下一节将深入解析 SQL 语句的构建与执行过程。

5. Spark内存计算与任务调度源码分析

Spark 是当前最主流的分布式内存计算框架之一,其核心优势在于将数据尽可能地保留在内存中,以提升大规模数据处理的性能。Spark 的执行模型、任务调度机制以及内存管理机制共同构成了其高效处理能力的基石。本章将深入分析 Spark 的执行模型、任务划分流程、内存管理机制以及调度系统,从源码层面揭示其高效处理背后的实现逻辑。

5.1 Spark执行模型与任务划分

Spark 的执行模型建立在 RDD(Resilient Distributed Dataset)基础之上,通过 DAG(有向无环图)将计算任务划分成多个阶段(Stage),每个阶段包含多个任务(Task),最终在集群的各个 Executor 上并行执行。

5.1.1 RDD的创建与DAG生成

RDD 是 Spark 的核心抽象,表示一个不可变、可分区、可并行计算的元素集合。用户可以通过读取外部数据源或转换现有 RDD 来创建新的 RDD。

RDD 创建示例:
val rdd = sc.textFile("hdfs://path/to/file")

该语句通过 SparkContext textFile 方法从 HDFS 读取文件并创建一个 RDD。 textFile 的源码如下(简化版):

def textFile(path: String, minPartitions: Int = defaultMinPartitions): RDD[String] = {
  hadoopFile(path, classOf[TextInputFormat], classOf[LongWritable], classOf[Text], minPartitions)
    .map(pair => pair._2.toString)
}
  • hadoopFile 方法基于 Hadoop 的 InputFormat 创建一个 RDD。
  • map 操作将 (LongWritable, Text) 转换为 String 类型。

当用户调用 RDD 的 map filter flatMap 等转换操作时,Spark 并不会立即执行,而是构建一个 DAG 来记录操作之间的依赖关系。

DAG 的构建过程:

Spark 使用 DAGScheduler 负责将 RDD 的转换操作转化为 DAG,并根据宽依赖(ShuffleDependency)将整个计算流程划分为多个 Stage。

// DAGScheduler.scala
def handleJobSubmitted(jobId: Int, ...) {
  val finalStage = createResultStage(...)
  submitStage(finalStage)
}

private def createResultStage(...) {
  val parents = getOrCreateParentStages(...)
  Stage(id, rdd, parentStages = parents, ...)
}
  • createResultStage 方法递归构建 Stage。
  • 每个 Stage 由一组具有窄依赖关系的 RDD 转换组成,遇到宽依赖则切分 Stage。

5.1.2 Stage划分与Task生成机制

Stage 划分的核心逻辑是基于 RDD 的依赖类型(窄依赖和宽依赖)进行切分。窄依赖表示父 RDD 的每个分区最多被子 RDD 的一个分区使用;宽依赖则表示父 RDD 的多个分区可能被子 RDD 的一个分区使用,此时需要触发 Shuffle。

Stage 划分流程图(mermaid):
graph TD
    A[用户提交Job] --> B{DAGScheduler创建FinalStage}
    B --> C[递归查找父Stage]
    C --> D[遇到ShuffleDependency则切分Stage]
    D --> E[继续向上查找]
    E --> F[构建Stage DAG]
    F --> G[将Stage拆分为TaskSet]
    G --> H[TaskScheduler提交任务到Executor]
Task 生成过程:

每个 Stage 最终会被拆分为多个 Task,Task 是 Spark 调度和执行的最小单位。Task 的生成主要由 TaskScheduler 完成。

// TaskSchedulerImpl.scala
override def submitTasks(taskSet: TaskSet) {
  val tasks = taskSet.tasks
  launchTasks(tasks)
}

private def launchTasks(tasks: Seq[Task[_]]) {
  for (task <- tasks) {
    execBackend.launchTask(task)
  }
}
  • launchTasks 方法将每个 Task 提交给 Executor 执行。
  • Task 包含了 RDD 分区的计算逻辑和执行上下文。

5.1.3 Task的执行与结果收集

Task 在 Executor 上执行时,会调用 RDD 的 compute 方法来计算每个分区的数据。

示例:Task 执行流程(代码简化):
// Executor.scala
def runTask(task: Task[_]) {
  val result = task.run()
  sendResult(result)
}
  • task.run() 方法会调用 RDD 的 compute 方法。
  • sendResult 将计算结果返回给 Driver。

结果收集由 DAGScheduler 负责,当所有 Task 成功执行后,Driver 会汇总结果并返回给用户。

// DAGScheduler.scala
def handleTaskCompletion(...) {
  if (allTasksCompleted(stage)) {
    val results = fetchTaskResults(...)
    job.listener.taskSucceeded(results)
  }
}
  • fetchTaskResults 收集所有 Task 的输出。
  • 最终通过 JobWaiter 等待结果并返回。

5.2 内存管理与执行引擎

Spark 的高性能很大程度上依赖其高效的内存管理机制。Spark 使用统一的内存模型,将内存划分为 ExecutionMemoryPool(执行内存)和 StorageMemoryPool(存储内存),并通过 Tungsten 引擎优化内存使用效率。

5.2.1 ExecutionMemoryPool与StorageMemoryPool的实现

Spark 的内存管理由 MemoryManager 负责,内存池分为两类:

内存池类型 用途说明
ExecutionMemoryPool 用于任务执行过程中的内存分配,如 Shuffle、Sort、Aggregation
StorageMemoryPool 用于缓存 RDD 数据和广播变量
内存池初始化流程:
// MemoryManager.scala
val executionMemoryPool = new ExecutionMemoryPool(totalMemory)
val storageMemoryPool = new StorageMemoryPool(totalMemory)
  • ExecutionMemoryPool 管理任务执行所需内存。
  • StorageMemoryPool 管理缓存所需内存。

Spark 通过动态调节机制在两者之间平衡内存使用,避免资源浪费。

5.2.2 Tungsten引擎的二进制存储优化

Tungsten 是 Spark 的二进制存储引擎,旨在减少 JVM 堆内存的使用和垃圾回收压力。它通过将数据以二进制格式存储在堆外内存中,提升数据处理效率。

Tungsten 内存结构示意图(mermaid):
graph LR
    A[BinaryFormat] --> B[UnsafeMemory]
    B --> C[MemoryManager]
    C --> D[ExecutionMemoryPool]
    C --> E[StorageMemoryPool]
  • 数据以二进制格式存储,节省内存开销。
  • 支持直接操作堆外内存,减少 GC 压力。
示例:使用 Tungsten 存储数据:
val buffer = new UnsafeRowWriter(1024)
buffer.writeInt(0, 42)
val row = buffer.getRow
  • UnsafeRowWriter 是 Tungsten 提供的写入器。
  • row 是一个二进制格式的 Row 对象。

5.2.3 Shuffle过程的内存使用与优化

Shuffle 是 Spark 中最耗内存的操作之一,涉及数据的分区、排序、序列化与传输。

Shuffle 过程内存使用分析表:
阶段 内存使用情况
Map 阶段 内存用于缓存输出数据,使用 MemoryBuffer
Shuffle 写阶段 序列化数据写入磁盘或内存,使用 ShuffleMemoryManager
Reduce 阶段 缓存拉取的数据,使用 AggregationBuffer
Shuffle 内存优化配置项:
配置项 默认值 说明
spark.shuffle.memoryFraction 0.2 用于 Shuffle 的堆内存比例
spark.storage.memoryFraction 0.6 用于存储的堆内存比例
spark.executor.memoryOverhead 384MB 堆外内存预留大小

通过合理配置这些参数,可以有效提升 Shuffle 性能并避免 OOM。

5.3 调度系统与资源分配

Spark 的调度系统由 DAGScheduler TaskScheduler 构成,负责将任务调度到合适的 Executor 上执行。Executor 的注册、任务分配策略以及动态资源分配机制共同构成了 Spark 的资源管理核心。

5.3.1 DAGScheduler与TaskScheduler的职责划分

模块 职责说明
DAGScheduler 将 RDD 转换为 DAG,划分 Stage,生成 TaskSet
TaskScheduler 接收 TaskSet,将 Task 分配给合适的 Executor 执行
DAGScheduler 与 TaskScheduler 协作流程图(mermaid):
graph LR
    A[用户提交Job] --> B[DAGScheduler划分Stage]
    B --> C[生成TaskSet]
    C --> D[TaskScheduler提交Task]
    D --> E[Executor执行Task]
    E --> F[结果返回Driver]
  • DAGScheduler 负责逻辑调度。
  • TaskScheduler 负责物理调度。

5.3.2 Executor的注册与任务分配策略

Executor 启动后会向 Driver 注册,并定期发送心跳信息以维持连接。

Executor 注册流程代码片段:
// Executor.scala
val driver = SparkEnv.get.driverEndpoint
driver.send(RegisterExecutor(...))
  • RegisterExecutor 消息通知 Driver 该 Executor 可用。

Driver 收到注册后,将其加入 Executor 列表,并开始分配任务。

Task 分配策略(Locality Level):
策略级别 说明
PROCESS_LOCAL Task 与数据在同一 Executor 上
NODE_LOCAL Task 与数据在同一节点上
NO_PREF 无偏好,随机分配

优先分配本地数据,减少网络传输开销。

5.3.3 动态资源分配的实现机制

Spark 支持动态资源分配,根据 Job 的负载自动申请或释放 Executor。

动态资源分配流程图(mermaid):
graph TD
    A[Driver启动] --> B{是否有空闲Executor?}
    B -->|是| C[直接分配任务]
    B -->|否| D[申请新Executor]
    D --> E[Executor启动并注册]
    E --> F[继续执行任务]
    F --> G[任务完成,释放Executor]
相关配置参数:
配置项 默认值 说明
spark.dynamicAllocation.enabled false 是否启用动态资源分配
spark.dynamicAllocation.maxExecutors 无穷大 最大 Executor 数量
spark.dynamicAllocation.idleTimeout 60s 空闲超时时间

启用动态资源分配可以有效节省资源,尤其适用于多用户共享集群的场景。

5.4 Spark Streaming流式处理

Spark Streaming 是 Spark 提供的实时流处理框架,采用微批处理模型(Micro-batch Processing)实现流式数据的实时处理。

5.4.1 DStream的构建与微批处理模型

DStream(Discretized Stream)是 Spark Streaming 的核心抽象,表示连续的数据流。

DStream 创建示例:
val lines = ssc.socketTextStream("localhost", 9999)
val words = lines.flatMap(_.split(" "))
val wordCounts = words.map(word => (word, 1)).reduceByKey(_ + _)
wordCounts.print()
  • socketTextStream 创建一个基于 Socket 的 DStream。
  • flatMap map 对数据进行转换。
  • reduceByKey 聚合数据。
  • print 输出结果。

微批处理模型将流式数据划分为小批次(Batch),每个批次在 Spark 中作为一个 RDD 进行处理。

5.4.2 Receiver机制与数据接收流程

Receiver 是 Spark Streaming 接收外部数据流的核心组件,运行在 Executor 上,负责持续接收数据并写入 BlockManager。

Receiver 数据接收流程图(mermaid):
graph LR
    A[数据源] --> B[Receiver接收数据]
    B --> C[写入BlockManager]
    C --> D[DAGScheduler调度处理]
    D --> E[结果输出]
Receiver 启动代码示例:
// ReceiverInputDStream.scala
def start() {
  val receiverExecutor = SparkEnv.get.executor
  receiverExecutor.submit(new Runnable {
    def run() {
      receiver.onStart()
      while (!isStopped) {
        val data = receiver.receive()
        store(data)
      }
    }
  })
}
  • onStart 启动接收逻辑。
  • receive 方法持续接收数据。
  • store 将数据写入 BlockManager。

5.4.3 Checkpoint机制与容错实现

为了实现流处理的容错性,Spark Streaming 提供了 Checkpoint 机制,定期将 DStream 的元数据和状态保存到可靠的存储系统中。

Checkpoint 配置示例:
ssc.checkpoint("hdfs://checkpoint-path")
  • 启用 Checkpoint 后,Spark 会定期保存 DStream 的状态。
Checkpoint 容错流程图(mermaid):
graph TD
    A[任务运行] --> B{是否到达Checkpoint间隔?}
    B -->|是| C[保存DStream状态]
    C --> D[HDFS持久化]
    D --> E[故障恢复时读取状态]
  • 若发生失败,Spark Streaming 会从最近的 Checkpoint 恢复状态,确保 Exactly-Once 语义。

本章内容完。

6. Flink流处理与状态管理源码分析

Flink 是当前流处理领域的佼佼者,其核心优势在于对状态的高效管理、低延迟的计算模型以及对 Exactly-Once 语义的天然支持。本章将深入剖析 Flink 的作业执行模型、状态管理机制、窗口与事件时间处理流程,以及底层网络通信与背压控制机制,帮助读者全面掌握其运行时架构与核心源码逻辑。

6.1 Flink作业执行模型与运行时架构

Flink 的作业执行模型围绕 JobGraph、ExecutionGraph、TaskManager 与 JobManager 的协同工作展开,构成了其分布式执行的核心。

6.1.1 JobGraph与ExecutionGraph的构建

Flink 作业的执行流程始于用户提交的 Job,Flink 会将其转换为 JobGraph ,这是作业的逻辑视图。JobGraph 包含多个 JobVertex,每个 JobVertex 表示一个操作符,如 Map、Filter、Window 等。

JobGraph 会被提交给 JobManager,JobManager 会将其进一步转换为 ExecutionGraph ,ExecutionGraph 是 JobGraph 的并行化版本,包含了每个 JobVertex 的多个 ExecutionVertex,每个 ExecutionVertex 对应一个 Task。

// 示例:JobGraph 构建过程(简化)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> source = env.addSource(new MySourceFunction());
DataStream<String> processed = source.map(new MyMapFunction());
processed.print();

JobGraph jobGraph = env.getStreamGraph().getJobGraph();

参数说明:
- StreamExecutionEnvironment :Flink 流处理的执行环境。
- addSource :添加数据源。
- map :转换操作。
- print :输出操作。
- getJobGraph() :将 StreamGraph 转换为 JobGraph。

6.1.2 TaskManager与JobManager的协作机制

JobManager 负责协调整个作业的执行,包括资源分配、任务调度、状态检查点等。

  • JobManager :负责作业的提交、ExecutionGraph 的生成、任务的调度。
  • TaskManager :负责执行具体的 Task,并与 JobManager 保持心跳通信。

在任务执行阶段,JobManager 会将 Task 分配给 TaskManager 执行。TaskManager 启动线程执行 Task,并将执行结果反馈给 JobManager。

6.1.3 Checkpoint与Savepoint的触发流程

Flink 的容错机制依赖于 Checkpoint 与 Savepoint。

  • Checkpoint :定期触发的全局状态快照,用于故障恢复。
  • Savepoint :手动触发的快照,用于升级、迁移等场景。

Checkpoint 流程如下:

  1. JobManager 向所有 Source Task 发送 Checkpoint Barrier。
  2. Source Task 接收到 Barrier 后,将当前状态快照写入状态后端(如 RocksDB)。
  3. 所有 Task 完成快照后,JobManager 接收通知,完成 Checkpoint。
graph TD
    A[JobManager] --> B[发送Checkpoint Barrier]
    B --> C[Source Task]
    C --> D[处理Barrier并写入状态]
    D --> E[下游Task依次处理Barrier]
    E --> F[所有Task完成快照]
    F --> G[JobManager确认Checkpoint完成]

下一节将深入探讨 Flink 的状态管理与容错机制。

本文还有配套的精品资源,点击获取 menu-r.4af5f7ec.gif

简介:源码分析是提升IT技术水平、深入理解软件运行机制的重要方法。本资源为名为“xmind-master”的压缩包,包含使用Xmind工具整理的多个主流Java及大数据框架的源码分析思维导图,涵盖Spring、SpringMVC、MyBatis、Spark、Flink及JVM等核心技术。通过这些思维导图,开发者可系统理解各框架的架构设计、核心流程、组件交互及设计模式,掌握并发控制、异常处理与JVM优化技巧,全面提升代码质量与问题解决能力。


本文还有配套的精品资源,点击获取
menu-r.4af5f7ec.gif

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐