从 MySQL 到 StarRocks:某电商平台报表采集与 Stream Load 实战
·
受众说明:本文面向熟悉 Java 服务与数据采集流程、对 StarRocks 有基本认知的读者。
某电商平台报表数据量激增后,MySQL 写入瓶颈明显,我们迁移到 StarRocks,并用 Stream Load 替代 JDBC 批量插入。迁移过程中遇到 FE→BE 重定向导致本地无法解析内网域名的问题,通过固定 FE 地址、配置 hosts 或手动处理重定向解决。表设计采用 PRIMARY KEY(unique_key),无需先删后增,Stream Load 自动覆盖更新。Buckets 选 12 以匹配集群并行度与读写平衡。
1. 背景:MySQL 已经扛不住
最初 报表落库流程很直接:
- 下载报表文档
- 解析 JSON
- 分批插入 MySQL
随着单次报表数据量达到百万级,MySQL 逐步暴露出问题:
- 写入耗时增长,任务执行时间不可控
- 批次越大越容易出现超时
- 数据采集整体吞吐被限制
结论很明确:MySQL 不再适合该场景,需要切换到分析型数据库。
2. 选型:为什么是 StarRocks
我们需要的能力是:
- 高吞吐写入
- 更适合分析型查询
- 可容忍一定延迟
StarRocks 在写入性能、OLAP 查询和批量导入方面表现更好,且 Stream Load 可以显著提升导入吞吐,因此确定迁移到 StarRocks。
3. 表结构设计:主键覆盖更新,省掉“先删后增”
关键思路是定义唯一主键,让导入行为变成“幂等覆盖”,从根源上消除“先删后增”的复杂度。
建表 SQL(核心字段)
CREATE TABLE amazon_sp_search_terms_xxx (
unique_key varchar(50) NOT NULL COMMENT "联合索引",
id bigint NOT NULL,
country_code varchar(10) NULL,
marketplace_id varchar(20) NULL,
search_term varchar(500) NULL,
search_frequency_rank varchar(50) NULL,
create_time bigint NOT NULL
)
ENGINE = OLAP
PRIMARY KEY(`unique_key`)
DISTRIBUTED BY HASH(`unique_key`) BUCKETS 12
PROPERTIES (
"replication_num" = "3"
);
为什么选 12 个 Buckets
- 并行度与集群规模平衡:12 通常能让 3~4 台 BE 形成良好并行度
- 写入与查询权衡:Buckets 太少会限制并发,太多会导致小文件碎片
- 便于后期扩展:12 是中等规模集群的合理起点,可结合监控再调优
4. Stream Load:百万级数据导入的核心能力
相比 JDBC 批量插入,Stream Load 的优势是:
- HTTP 方式导入,吞吐更高
- 支持并发写入
- 可以直接用 CSV/TSV 文件落库
基本请求格式
curl -u : -H “format: csv”
-H “column_separator:\t”
-H “columns: unique_key,id,country_code,marketplace_id,search_term,search_frequency_rank,create_time”
-T data.tsv
“http://fe_host:8030/api/db/xxx/_stream_load”
Java 参考实现(可直接运行的 Stream Load 测试)
package com.test.dcs.data.amazon;
import com.test.common.exception.testValidateException;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.commons.lang3.StringUtils;
import org.apache.http.Header;
import org.apache.http.HttpEntity;
import org.apache.http.HttpHeaders;
import org.apache.http.client.config.RequestConfig;
import org.apache.http.client.methods.CloseableHttpResponse;
import org.apache.http.client.methods.HttpPut;
import org.apache.http.entity.StringEntity;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClients;
import org.apache.http.util.EntityUtils;
import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.Base64;
import java.util.Map;
import java.util.UUID;
/**
* @Description: Stream Load 单元测试类
* @Author: azu
* @CreateTime: 2026-03-03 15:45
* @Version: 1.0
*/
public class AmazonSpSearchTermsStreamLoadTest {
private static final Logger logger = LoggerFactory.getLogger(AmazonSpSearchTermsStreamLoadTest.class);
private static final String DEFAULT_STREAM_LOAD_URL = "http://fe_host:8030/api/{db_name}/{table_name}/_stream_load";
private static final String DEFAULT_USERNAME = "your_user";
private static final String DEFAULT_PASSWORD = "your_password";
private static final String DEFAULT_CSV_PATH = "C:\\path\\to\\sp_search_terms_stream_test.csv";
private static final String DEFAULT_COLUMNS = "unique_key,id,country_code,marketplace_id,xxx
/**
* @Description: Stream Load CSV 文件到 StarRocks
* @Author: azu
* @CreateTime: 2026-03-03 15:45
*/
@Test
public void testStreamLoadFromCsv() throws Exception {
String streamLoadUrl = resolveParam("sr.stream.url", DEFAULT_STREAM_LOAD_URL);
String username = resolveParam("sr.user", DEFAULT_USERNAME);
String password = resolveParam("sr.pass", DEFAULT_PASSWORD);
String csvPath = resolveParam("sr.csv", DEFAULT_CSV_PATH);
String columns = resolveParam("sr.columns", DEFAULT_COLUMNS);
boolean followRedirect = Boolean.parseBoolean(resolveParam("sr.follow.redirect", "false"));
boolean allowInternalRedirect = Boolean.parseBoolean(resolveParam("sr.allow.internal.redirect", "false"));
validateRequired(streamLoadUrl, "Stream Load URL 不能为空");
validateRequired(username, "Stream Load 用户名不能为空");
validateRequired(password, "Stream Load 密码不能为空");
validateRequired(csvPath, "CSV 文件路径不能为空");
String content = new String(Files.readAllBytes(Paths.get(csvPath)), StandardCharsets.UTF_8);
RequestConfig requestConfig = RequestConfig.custom()
.setConnectTimeout(60 * 1000)
.setSocketTimeout(30 * 60 * 1000)
.setRedirectsEnabled(false)
.build();
try (CloseableHttpClient httpClient = HttpClients.custom().setDefaultRequestConfig(requestConfig).build()) {
String auth = Base64.getEncoder().encodeToString((username + ":" + password).getBytes(StandardCharsets.UTF_8));
String label = "search_terms_" + System.currentTimeMillis() + "_" + UUID.randomUUID().toString().replace("-", "");
StreamLoadResult result = executeStreamLoad(httpClient, streamLoadUrl, auth, label, columns, content);
if (isRedirectStatus(result.getStatusCode())) {
if (StringUtils.isBlank(result.getLocation())) {
logger.error("Stream Load 重定向,状态码:{},Location:{}", result.getStatusCode(), result.getLocation());
throw new testValidateException("Stream Load 重定向失败:" + result.getLocation());
}
if (followRedirect && (allowInternalRedirect || !isInternalHost(result.getLocation()))) {
result = executeStreamLoad(httpClient, result.getLocation(), auth, label, columns, content);
} else {
logger.error("Stream Load 重定向,状态码:{},Location:{}", result.getStatusCode(), result.getLocation());
throw new testValidateException("Stream Load 重定向失败:" + result.getLocation());
}
}
if (result.getStatusCode() != 200 || !isStreamLoadSuccess(result.getBody())) {
logger.error("Stream Load 失败,状态码:{},响应体:{}", result.getStatusCode(), result.getBody());
throw new testValidateException("Stream Load 失败:" + result.getBody());
}
}
}
private StreamLoadResult executeStreamLoad(CloseableHttpClient httpClient, String url, String auth, String label, String columns, String content) throws Exception {
HttpPut put = new HttpPut(url);
put.setHeader(HttpHeaders.AUTHORIZATION, "Basic " + auth);
put.setHeader("label", label);
put.setHeader("format", "csv");
put.setHeader("column_separator", "\t");
put.setHeader("columns", columns);
put.setHeader("timeout", "1800");
put.setHeader(HttpHeaders.EXPECT, "100-continue");
put.setHeader("Content-Type", "text/plain; charset=UTF-8");
put.removeHeaders(HttpHeaders.CONTENT_LENGTH);
put.setEntity(new StringEntity(content, StandardCharsets.UTF_8));
try (CloseableHttpResponse response = httpClient.execute(put)) {
HttpEntity entity = response.getEntity();
String body = entity == null ? "" : EntityUtils.toString(entity, StandardCharsets.UTF_8);
Header location = response.getFirstHeader(HttpHeaders.LOCATION);
String locationValue = location == null ? "" : location.getValue();
return new StreamLoadResult(response.getStatusLine().getStatusCode(), body, locationValue);
}
}
private boolean isRedirectStatus(int statusCode) {
return statusCode == 301 || statusCode == 302 || statusCode == 307 || statusCode == 308;
}
private boolean isInternalHost(String url) {
if (StringUtils.isBlank(url)) {
return false;
}
String lower = url.toLowerCase();
return lower.contains(".svc.cluster.local") || lower.contains("starrocks-be");
}
private boolean isStreamLoadSuccess(String responseBody) throws Exception {
if (StringUtils.isBlank(responseBody)) {
return false;
}
ObjectMapper mapper = new ObjectMapper();
Map<String, Object> data = mapper.readValue(responseBody, Map.class);
Object status = data.get("Status");
if (status == null) {
status = data.get("status");
}
return "Success".equalsIgnoreCase(String.valueOf(status));
}
private String resolveParam(String key, String defaultValue) {
String value = System.getProperty(key);
if (StringUtils.isBlank(value)) {
value = System.getenv(key);
}
return StringUtils.isBlank(value) ? defaultValue : value;
}
private void validateRequired(String value, String message) {
if (StringUtils.isBlank(value)) {
throw new testValidateException(message);
}
}
private static class StreamLoadResult {
private final int statusCode;
private final String body;
private final String location;
private StreamLoadResult(int statusCode, String body, String location) {
this.statusCode = statusCode;
this.body = body;
this.location = location;
}
private int getStatusCode() {
return statusCode;
}
private String getBody() {
return body;
}
private String getLocation() {
return location;
}
}
}
关键点
- 不需要删除旧数据
主键一致会自动覆盖旧行,导入过程天然幂等 - 适合超大批量
适用于单次百万级数据的导入场景
5. 踩坑:FE 重定向到 BE,本地无法解析内网域名
StarRocks 的 Stream Load 经常返回 307 重定向,FE 会将请求重定向到 BE。
如果本地无法解析 BE 内网域名,会报错:
getaddrinfo ENOTFOUND starrocks-be-*.svc.cluster.local
解决方式
-
优先直连 FE 地址
使用可访问的 FE 地址与端口,避免网关或入口产生二次重定向 -
配置 hosts 映射
将 BE 域名映射到可访问 IP,例如:
10.x.x.x starrocks-be-0.xxx.local
10.x.x.x starrocks-be-1.xxx.local
10.x.x.x starrocks-be-2.xxx.local
- 应用层处理重定向
如果收到 307/302,读取 Location 并重新发起请求
6. 迁移后收益
- 写入速度显著提升
- 导入稳定性更好
- 逻辑简化,去掉先删后增
- 报表数据落库可支撑更大规模
7. 总结
这次迁移的核心变化是两点:
- MySQL → StarRocks,解决写入瓶颈
- JDBC 批量 → Stream Load,提升吞吐与稳定性
配合 PRIMARY KEY 的幂等更新机制,整个落库逻辑更清晰、更稳定。
如果后续数据量继续增长,可继续从 buckets、压缩、并发导入策略等方面做进一步调优。
PS:这篇文章使用technical-blog-writing skills来编写,比起以往的编辑方式,现在真是越来越高效了
更多推荐



所有评论(0)