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 协议下载,可并发、可指定行范围 |
两个关键事实(封装必须依赖):
Odps对象线程安全,可作为单例 Bean 全局复用(new Odps(new AliyunAccount(ak, sk))一次,处处共享)。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. 可行性验证步骤
完成编码后,按以下步骤验证(缺一不可):
- 连通性验证:
Odps odps = OdpsClientFactory.create(props);后调用odps.projects().exists(project)返回true。 - 小结果集:
SELECT 1 AS a, 'hello' AS b应返回[{a=1, b=hello}]。 - 列名正确性:
SELECT user_id, name FROM t LIMIT 5的 Map key 与列名一致。 - 超时生效:故意查一张超大表,确认在
queryTimeoutSeconds后抛MaxComputeQueryTimeoutException而非无限等待。 - 大结果集:构造 >10000 行的查询,走
queryStream能完整读回且不 OOM。 - 注入拦截:提交
SELECT * FROM t; DROP TABLE t,确认被SqlValidator拒绝。 - 并发安全:多线程同时调用
query(),确认单例 Odps 无异常。
11. 测试方案
| 层级 | 内容 | 工具 |
|---|---|---|
| 单元测试 | SqlValidator 白名单/黑名单用例、RecordMapper 类型映射用例 |
JUnit 5 |
| 集成测试 | 连真实测试项目跑 SELECT 1、超时、大结果集 |
需测试项目 + 测试 AK |
| 性能测试 | 并发查询压测,观察线程池、QPS、延迟 | JMH / 压测脚本 |
12. 性能优化
- 单例复用 Odps(已做),避免重复建连接。
- 大结果集走流式,避免一次性
List<Record>撑爆内存。 - 内网 endpoint:ECS 上务必用内网 endpoint,省流量、低延迟。
- 分页/裁剪:查询层强制
maxRows+ 建议业务侧加LIMIT,避免全表扫描。 - 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> + AutoCloseable:hasNext() / 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 以控制台「项目信息」页面为准。