MaxCompute 数据查询 SDK 封装 — 技术方案

版本:v1.0 日期:2026-08-13 适用环境:Java 11 / Spring Boot 2.7 / Maven / 阿里云 MaxCompute 目标模块:cx-health 项目的 cx-third-util(第三方集成模块)


1. 背景与目标

为 Spring Cloud 微服务提供统一、类型安全、可控超时、防注入的 MaxCompute 数据查询能力。

核心目标:

目标 说明
统一入口 业务代码只调用一个 MaxComputeQueryService,不接触 SDK 细节
类型安全 自动完成 MaxCompute 类型 → Java 类型的映射
可控超时 查询超时、结果上限均可配置,避免线程堆积
大数据量 小结果集走内存读取,大结果集走流式 / Tunnel 下载
安全合规 SQL 白名单校验 + 权限最小化 + 审计日志

2. 技术选型

选择 版本 说明
核心 SDK com.aliyun.odps:odps-sdk-core 0.59.0-public 官方 Java SDK,含 Odps 客户端、SQLTask、InstanceTunnel
传递依赖 odps-sdk-commons 随 core 传递 含 Record、TableSchema、Column、TypeInfo 等数据模型
认证 RAM 子账号 AccessKey - 生产建议 RAM Role + STS 临时凭证
Java 11 - 与 cx-health 一致

Maven 依赖(odps-sdk-core 会自动传递 odps-sdk-commons,无需手动加):

<dependency>
    <groupId>com.aliyun.odps</groupId>
    <artifactId>odps-sdk-core</artifactId>
    <version>0.59.0-public</version>
</dependency>

仓库地址:阿里云 maven 中央仓库镜像 https://maven.aliyun.com/repository/public(国内网络必须配,否则拉不动国外中央仓库)。


3. 核心原理(关键认知)

MaxCompute 查询是异步作业模型,不是 JDBC 那种同步 executeQuery

提交 SQL → 得到 Instance(作业) → 轮询等待完成 → 读取结果

三条读取路径,按结果集规模选择:

场景 API 特点
小结果集(≤ 10000 行) SQLTask.getResult(instance) 一次性返回 List<Record>,简单但内存敏感
大结果集(> 10000 行) SQLTask.getResultSet(instance) 流式 ResultSet(实现 Iterator + AutoCloseable),内存恒定
超大结果 / 需断点续传 InstanceTunnel.createDownloadSession(...) 走 Tunnel 协议下载,可并发、可指定行范围

两个关键事实(封装必须依赖):

  1. Odps 对象线程安全,可作为单例 Bean 全局复用(new Odps(new AliyunAccount(ak, sk)) 一次,处处共享)。
  2. SQLTask.getResult() 默认上限 10000 行,超出会被截断——大表查询必须显式走 getResultSet 或 Tunnel。

4. 整体架构

┌─────────────────────────────────────────────────────┐
│  Controller / 业务 Service                            │
└───────────────────────┬─────────────────────────────┘
                        ▼
┌─────────────────────────────────────────────────────┐
│  MaxComputeQueryService(执行层,对外唯一入口)        │
│   - query(sql)          → List<Map>  小结果集         │
│   - queryStream(sql)    → 流式回调   大结果集         │
│   - query(sql, limit)   → 带行数上限                  │
│   内含:超时、线程池、限流、审计、SQL 校验             │
└───────────────────────┬─────────────────────────────┘
                        ▼
┌─────────────────────────────────────────────────────┐
│  OdpsClientFactory(客户端层,单例 Odps Bean)         │
└───────────────────────┬─────────────────────────────┘
                        ▼
┌─────────────────────────────────────────────────────┐
│  ResultMapper(结果映射层)Record → Map / Java 对象    │
└─────────────────────────────────────────────────────┘

模块目录结构(建议落在 cx-third-util 下):

com.example.maxcompute
├── config/
│   ├── MaxComputeProperties.java         # 配置属性
│   └── MaxComputeAutoConfiguration.java  # 自动装配
├── client/
│   └── OdpsClientFactory.java            # 客户端工厂
├── exception/
│   ├── MaxComputeException.java          # 业务异常
│   └── MaxComputeQueryTimeoutException.java
├── mapper/
│   └── RecordMapper.java                 # 结果映射
├── service/
│   ├── MaxComputeQueryService.java       # 核心查询服务
│   └── QueryAuditLogger.java             # 审计
└── support/
    └── SqlValidator.java                 # SQL 白名单校验

5. 详细设计

5.1 配置层

package com.example.maxcompute.config;

import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;

@Data
@ConfigurationProperties(prefix = "maxcompute")
public class MaxComputeProperties {

    /** 内网/公网 endpoint(必填) */
    private String endpoint;

    /** Tunnel 下载 endpoint(可选,大数据量用) */
    private String tunnelEndpoint;

    /** 项目名(必填) */
    private String project;

    /** AccessKey ID(必填,或走 STS) */
    private String accessKeyId;

    /** AccessKey Secret(必填,或走 STS) */
    private String accessKeySecret;

    /** STS 安全令牌(可选,使用 RAM Role 时填) */
    private String securityToken;

    /** 查询超时(秒),默认 300 */
    private long queryTimeoutSeconds = 300;

    /** 结果行数上限,默认 10000 */
    private long maxRows = 10000;

    /** 查询线程池核心线程数 */
    private int poolCoreSize = 8;

    /** 查询线程池最大线程数 */
    private int poolMaxSize = 16;

    /** 线程池队列容量 */
    private int poolQueueCapacity = 200;
}

application.yml

maxcompute:
  endpoint: http://service.cn-hangzhou.maxcompute.aliyun.com/api
  tunnel-endpoint: http://dt.cn-hangzhou.maxcompute.aliyun.com
  project: your_project
  access-key-id: ${MAXCOMPUTE_AK}
  access-key-secret: ${MAXCOMPUTE_SK}
  # security-token: ${MAXCOMPUTE_STS_TOKEN}
  query-timeout-seconds: 300
  max-rows: 10000
  pool-core-size: 8
  pool-max-size: 16
  pool-queue-capacity: 200

5.2 客户端工厂(单例 Odps)

package com.example.maxcompute.client;

import com.aliyun.odps.Odps;
import com.aliyun.odps.account.Account;
import com.aliyun.odps.account.AliyunAccount;
import com.aliyun.odps.account.StsAccount;
import com.example.maxcompute.config.MaxComputeProperties;
import org.springframework.util.StringUtils;

public final class OdpsClientFactory {

    private OdpsClientFactory() {}

    /**
     * Odps 线程安全,作为单例 Bean 复用。
     * 有 securityToken 时走 STS 临时凭证,否则走 AccessKey。
     */
    public static Odps create(MaxComputeProperties props) {
        Account account;
        if (StringUtils.hasText(props.getSecurityToken())) {
            account = new StsAccount(props.getAccessKeyId(),
                                     props.getAccessKeySecret(),
                                     props.getSecurityToken());
        } else {
            account = new AliyunAccount(props.getAccessKeyId(),
                                        props.getAccessKeySecret());
        }
        Odps odps = new Odps(account);
        odps.setEndpoint(props.getEndpoint());
        odps.setDefaultProject(props.getProject());
        if (StringUtils.hasText(props.getTunnelEndpoint())) {
            odps.setTunnelEndpoint(props.getTunnelEndpoint());
        }
        return odps;
    }
}

5.3 异常体系

package com.example.maxcompute.exception;

/** 统一业务异常 */
public class MaxComputeException extends RuntimeException {
    public MaxComputeException(String message) { super(message); }
    public MaxComputeException(String message, Throwable cause) { super(message, cause); }
}
package com.example.maxcompute.exception;

/** 查询超时异常 */
public class MaxComputeQueryTimeoutException extends MaxComputeException {
    public MaxComputeQueryTimeoutException(String message, Throwable cause) {
        super(message, cause);
    }
}

5.4 结果映射层(类型转换)

MaxCompute 类型 → Java 类型映射:

MaxCompute Java 说明
STRING String 直接映射
BIGINT Long getBigint()
DOUBLE Double getDouble()
DECIMAL BigDecimal getDecimal()
DATETIME java.util.Date getDatetime()
BOOLEAN Boolean getBoolean()
ARRAY / MAP / STRUCT Object(原始或 JSON 序列化) 复合类型按需处理
package com.example.maxcompute.mapper;

import com.aliyun.odps.Column;
import com.aliyun.odps.OdpsType;
import com.aliyun.odps.TableSchema;
import com.aliyun.odps.data.Record;
import org.springframework.stereotype.Component;

import java.util.LinkedHashMap;
import java.util.Map;

@Component
public class RecordMapper {

    /** Record → Map(列名小写保持一致性可选) */
    public Map<String, Object> toMap(Record record, TableSchema schema) {
        Map<String, Object> row = new LinkedHashMap<>();
        for (int i = 0; i < schema.getColumns().size(); i++) {
            Column col = schema.getColumn(i);
            row.put(col.getName(), toJavaValue(record.get(i), col));
        }
        return row;
    }

    private Object toJavaValue(Object value, Column column) {
        if (value == null) {
            return null;
        }
        OdpsType type = column.getTypeInfo().getOdpsType();
        switch (type) {
            case DECIMAL:
                return value instanceof java.math.BigDecimal
                        ? value : new java.math.BigDecimal(value.toString());
            case DATETIME:
                return value; // SDK 已返回 java.util.Date
            case ARRAY:
            case MAP:
            case STRUCT:
                // 复合类型:默认返回原始对象,可在此接入 JSON 序列化
                return value;
            default:
                return value;
        }
    }
}

5.5 SQL 校验(防注入 + 防危险语句)

package com.example.maxcompute.support;

import org.springframework.stereotype.Component;

import java.util.Arrays;
import java.util.List;

@Component
public class SqlValidator {

    /** 禁止出现的关键字(即使包含在字符串里也拒绝,宁可误杀) */
    private static final List<String> FORBIDDEN = Arrays.asList(
            "DROP", "DELETE", "UPDATE", "INSERT", "CREATE",
            "ALTER", "TRUNCATE", "GRANT", "REVOKE", "MERGE");

    public void validate(String sql) {
        if (sql == null || sql.trim().isEmpty()) {
            throw new IllegalArgumentException("SQL 不能为空");
        }
        String upper = sql.trim().toUpperCase();
        // 仅允许查询类语句
        if (!(upper.startsWith("SELECT") || upper.startsWith("WITH"))) {
            throw new IllegalArgumentException("仅允许 SELECT / WITH 查询语句");
        }
        // 禁止危险关键字
        for (String kw : FORBIDDEN) {
            if (upper.contains(kw)) {
                throw new IllegalArgumentException("SQL 包含禁止关键字: " + kw);
            }
        }
    }
}

5.6 核心查询服务

package com.example.maxcompute.service;

import com.aliyun.odps.Instance;
import com.aliyun.odps.Odps;
import com.aliyun.odps.OdpsException;
import com.aliyun.odps.TableSchema;
import com.aliyun.odps.data.Record;
import com.aliyun.odps.data.ResultSet;
import com.aliyun.odps.task.SQLTask;
import com.example.maxcompute.config.MaxComputeProperties;
import com.example.maxcompute.exception.MaxComputeException;
import com.example.maxcompute.exception.MaxComputeQueryTimeoutException;
import com.example.maxcompute.mapper.RecordMapper;
import com.example.maxcompute.support.SqlValidator;
import org.springframework.stereotype.Service;

import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.function.Consumer;

@Service
public class MaxComputeQueryService {

    private final Odps odps;
    private final MaxComputeProperties props;
    private final RecordMapper recordMapper;
    private final SqlValidator sqlValidator;
    private final QueryAuditLogger auditLogger;

    public MaxComputeQueryService(Odps odps,
                                  MaxComputeProperties props,
                                  RecordMapper recordMapper,
                                  SqlValidator sqlValidator,
                                  QueryAuditLogger auditLogger) {
        this.odps = odps;
        this.props = props;
        this.recordMapper = recordMapper;
        this.sqlValidator = sqlValidator;
        this.auditLogger = auditLogger;
    }

    /** 同步查询(小结果集) */
    public List<Map<String, Object>> query(String sql) {
        long start = System.currentTimeMillis();
        Instance instance = submit(sql);
        waitForSuccess(instance);

        List<Record> records;
        try {
            records = SQLTask.getResult(instance);
        } catch (OdpsException e) {
            throw new MaxComputeException("读取查询结果失败", e);
        }
        TableSchema schema = getSchema(instance);
        List<Map<String, Object>> result = new ArrayList<>(records.size());
        for (Record r : records) {
            result.add(recordMapper.toMap(r, schema));
        }
        auditLogger.log(sql, result.size(), System.currentTimeMillis() - start);
        return result;
    }

    /** 流式查询(大结果集,内存友好) */
    public void queryStream(String sql, Consumer<Map<String, Object>> consumer) {
        long start = System.currentTimeMillis();
        Instance instance = submit(sql);
        waitForSuccess(instance);

        long count = 0;
        try (ResultSet rs = SQLTask.getResultSet(instance)) {
            TableSchema schema = rs.getTableSchema();
            while (rs.hasNext()) {
                consumer.accept(recordMapper.toMap(rs.next(), schema));
                count++;
            }
        } catch (Exception e) {
            throw new MaxComputeException("流式读取结果失败", e);
        }
        auditLogger.log(sql, count, System.currentTimeMillis() - start);
    }

    /** 带行数上限的查询(内部包一层 LIMIT) */
    public List<Map<String, Object>> query(String sql, long limit) {
        return query("SELECT * FROM (" + sql + ") __t LIMIT " + limit);
    }

    /** 提交 SQL 任务 */
    private Instance submit(String sql) {
        sqlValidator.validate(sql);
        try {
            return SQLTask.run(odps, sql);
        } catch (OdpsException e) {
            throw new MaxComputeException("提交 SQL 失败: " + sql, e);
        }
    }

    /** 阻塞等待完成,带超时 */
    private void waitForSuccess(Instance instance) {
        try {
            instance.waitForSuccess(props.getQueryTimeoutSeconds() * 1000);
        } catch (OdpsException e) {
            throw new MaxComputeQueryTimeoutException(
                    "查询超时(" + props.getQueryTimeoutSeconds() + "s),instanceId="
                            + instance.getId(), e);
        }
    }

    /** 获取结果 schema */
    private TableSchema getSchema(Instance instance) {
        try {
            return SQLTask.getResultSchema(instance);
        } catch (OdpsException e) {
            throw new MaxComputeException("获取结果 schema 失败", e);
        }
    }
}

5.7 审计日志

package com.example.maxcompute.service;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;

@Component
public class QueryAuditLogger {

    private static final Logger log = LoggerFactory.getLogger(QueryAuditLogger.class);

    /** 记录每次查询的 SQL、行数、耗时(可扩展落数据库/日志表) */
    public void log(String sql, long rows, long costMs) {
        log.info("[MaxCompute] rows={}, costMs={}, sql={}", rows, costMs, sql);
    }
}

5.8 自动装配

package com.example.maxcompute.config;

import com.aliyun.odps.Odps;
import com.example.maxcompute.client.OdpsClientFactory;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
@EnableConfigurationProperties(MaxComputeProperties.class)
public class MaxComputeAutoConfiguration {

    /** Odps 单例 Bean(线程安全) */
    @Bean
    public Odps odps(MaxComputeProperties props) {
        return OdpsClientFactory.create(props);
    }
}

6. 大数据量方案(> 10000 行)

方案 A(推荐):流式 ResultSet

// 内存恒定,适合百万级行数
maxComputeQueryService.queryStream(
    "SELECT user_id, amount FROM large_table WHERE dt='2026-08-01'",
    row -> process(row)   // 逐行处理,避免一次性加载全部
);

ResultSet 实现了 Iterator<Record> + AutoCloseable,可直接用于 try-with-resources 或增强 for 循环。

方案 B:结果落临时表 + Tunnel 下载

适合需要断点续传、并发下载、超大结果落盘:

import com.aliyun.odps.Instance;
import com.aliyun.odps.task.SQLTask;
import com.aliyun.odps.tunnel.InstanceTunnel;
import com.aliyun.odps.tunnel.io.TunnelRecordReader;
import com.aliyun.odps.data.Record;

// 1. 执行 SQL 并把结果写入带生命周期的临时表
String tmpTable = "tmp_query_" + UUID.randomUUID().toString().replace("-", "");
SQLTask.run(odps, "CREATE TABLE " + tmpTable + " LIFECYCLE 1 AS " + sql)
        .waitForSuccess();

// 2. 定位写结果的 instance,用 InstanceTunnel 下载
InstanceTunnel tunnel = new InstanceTunnel(odps);
InstanceTunnel.DownloadSession session =
        tunnel.createDownloadSession(project, instanceId);
try (TunnelRecordReader reader = session.openRecordReader(0, count)) {
    Record record;
    while ((record = reader.read()) != null) {
        // 处理 record
    }
}

LIFECYCLE 1 让临时表 1 天后自动清理,避免残留。


7. Spring Boot 集成

spring.factories(Spring Boot 2.7)或 AutoConfiguration.imports(2.7+ 用 META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports):

com.example.maxcompute.config.MaxComputeAutoConfiguration

业务侧用法:

@RestController
@RequestMapping("/api/maxcompute")
public class MaxComputeController {

    private final MaxComputeQueryService service;

    public MaxComputeController(MaxComputeQueryService service) {
        this.service = service;
    }

    @PostMapping("/query")
    public List<Map<String, Object>> query(@RequestBody Map<String, String> body) {
        return service.query(body.get("sql"));
    }
}

8. 认证方式

方式 适用场景 说明
AccessKey(AK/SK) 本地开发、非阿里云环境 new AliyunAccount(ak, sk)
STS 临时凭证 部署在 ECS / 有 RAM Role new StsAccount(ak, sk, token)new AliyunAccount(ak, sk, token)
RAM Role ECS 挂角色,免 AK 通过 STS SDK 换取临时凭证后走 StsAccount

权限最小化建议: RAM 子账号只授 odps:Read / odps:Select(数据读取权限),不给建表/删表等写权限。


9. 安全与成本

措施
权限最小化 RAM 只读策略,生产走 STS 临时凭证
SQL 注入 SqlValidator 只允许 SELECT/WITH + 危险关键字拦截
成本控制 MaxCompute 按扫描量计费:设置 maxRows 上限 + 单用户 QPS 限流
超时保护 waitForSuccess(timeout) + 独立线程池,防长查询打满线程
审计 记录 SQL、行数、耗时(可扩展落日志表)
凭据安全 AccessKey 放环境变量 / 配置中心,禁止硬编码进代码仓库

线程池配置(防止阻塞查询打满容器线程): 查询是阻塞式(waitForSuccess),必须用独立线程池承载,不要用 Tomcat 工作线程直接同步等待。建议通过 @Async 或自定义线程池把查询扔到 poolCoreSize=8 / poolMaxSize=16 / queueCapacity=200 的池子里。


10. 可行性验证步骤

完成编码后,按以下步骤验证(缺一不可):

  1. 连通性验证Odps odps = OdpsClientFactory.create(props); 后调用 odps.projects().exists(project) 返回 true
  2. 小结果集SELECT 1 AS a, 'hello' AS b 应返回 [{a=1, b=hello}]
  3. 列名正确性SELECT user_id, name FROM t LIMIT 5 的 Map key 与列名一致。
  4. 超时生效:故意查一张超大表,确认在 queryTimeoutSeconds 后抛 MaxComputeQueryTimeoutException 而非无限等待。
  5. 大结果集:构造 >10000 行的查询,走 queryStream 能完整读回且不 OOM。
  6. 注入拦截:提交 SELECT * FROM t; DROP TABLE t,确认被 SqlValidator 拒绝。
  7. 并发安全:多线程同时调用 query(),确认单例 Odps 无异常。

11. 测试方案

层级 内容 工具
单元测试 SqlValidator 白名单/黑名单用例、RecordMapper 类型映射用例 JUnit 5
集成测试 连真实测试项目跑 SELECT 1、超时、大结果集 需测试项目 + 测试 AK
性能测试 并发查询压测,观察线程池、QPS、延迟 JMH / 压测脚本

12. 性能优化

  1. 单例复用 Odps(已做),避免重复建连接。
  2. 大结果集走流式,避免一次性 List<Record> 撑爆内存。
  3. 内网 endpoint:ECS 上务必用内网 endpoint,省流量、低延迟。
  4. 分页/裁剪:查询层强制 maxRows + 建议业务侧加 LIMIT,避免全表扫描。
  5. SQA(可选):新版 SDK 支持交互式分析(低延迟即席查询),对延迟敏感的场景可后续评估接入,当前经典 SQLTask 流程最稳。

13. 常见坑与规避

现象 规避
getResult 截断 大表只返回 10000 行 显式走 getResultSet 或 Tunnel
异步模型误解 误以为 SQLTask.run 立即返回结果 必须 waitForSuccess 后再读结果
线程阻塞 同步等待打满 Tomcat 线程 独立线程池 + 超时
Endpoint 选错 公网 endpoint 慢/不通 ECS 内网用内网 endpoint
复合类型 ARRAY/MAP/STRUCT 返回原始对象难消费 映射层 JSON 序列化
AK 泄漏 密钥硬编码进 git 环境变量 / 配置中心

附录 A:核心 API 参考(已核实 0.59.0-public)

关键方法
Odps new Odps(Account) / setEndpoint / setDefaultProject / setTunnelEndpoint
Account 实现类 AliyunAccount(ak, sk[, token])StsAccount(ak, sk, token)
SQLTask run(Odps, sql) → Instance;getResult(Instance) → List\<Record>;getResultSet(Instance) → ResultSet;getResultSchema(Instance) → TableSchema
Instance waitForSuccess() / waitForSuccess(long ms) / getId() / isTerminated() / getTaskStatus()
ResultSet 实现 Iterator<Record> + AutoCloseablehasNext() / next() / getTableSchema() / getRecordCount()
Record get(int/String) / getBigint / getDouble / getBoolean / getDatetime / getDecimal / getColumnCount() / getColumns()
TableSchema getColumns() → List\<Column> / getColumn(int/String)
Column getName() / getTypeInfo() / getComment()
TypeInfo getOdpsType() → OdpsType / getTypeName()
InstanceTunnel createDownloadSession(project, instanceId) → DownloadSession

附录 B:Endpoint 参考

Region 公网 endpoint 内网 endpoint
华东1 杭州 http://service.cn-hangzhou.maxcompute.aliyun.com/api http://service.cn-hangzhou.maxcompute.aliyun-inc.com/api
华东2 上海 http://service.cn-shanghai.maxcompute.aliyun.com/api http://service.cn-shanghai.maxcompute.aliyun-inc.com/api
华北2 北京 http://service.cn-beijing.maxcompute.aliyun.com/api http://service.cn-beijing.maxcompute.aliyun-inc.com/api
华南1 深圳 http://service.cn-shenzhen.maxcompute.aliyun.com/api http://service.cn-shenzhen.maxcompute.aliyun-inc.com/api

准确 endpoint 以控制台「项目信息」页面为准。