受众说明:本文面向熟悉 Java 服务与数据采集流程、对 StarRocks 有基本认知的读者。

某电商平台报表数据量激增后,MySQL 写入瓶颈明显,我们迁移到 StarRocks,并用 Stream Load 替代 JDBC 批量插入。迁移过程中遇到 FE→BE 重定向导致本地无法解析内网域名的问题,通过固定 FE 地址、配置 hosts 或手动处理重定向解决。表设计采用 PRIMARY KEY(unique_key),无需先删后增,Stream Load 自动覆盖更新。Buckets 选 12 以匹配集群并行度与读写平衡。


1. 背景:MySQL 已经扛不住

最初 报表落库流程很直接:

  1. 下载报表文档
  2. 解析 JSON
  3. 分批插入 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

解决方式

  1. 优先直连 FE 地址
    使用可访问的 FE 地址与端口,避免网关或入口产生二次重定向

  2. 配置 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
  1. 应用层处理重定向
    如果收到 307/302,读取 Location 并重新发起请求

6. 迁移后收益

  • 写入速度显著提升
  • 导入稳定性更好
  • 逻辑简化,去掉先删后增
  • 报表数据落库可支撑更大规模

7. 总结

这次迁移的核心变化是两点:

  1. MySQL → StarRocks,解决写入瓶颈
  2. JDBC 批量 → Stream Load,提升吞吐与稳定性

配合 PRIMARY KEY 的幂等更新机制,整个落库逻辑更清晰、更稳定。
如果后续数据量继续增长,可继续从 buckets、压缩、并发导入策略等方面做进一步调优。

PS:这篇文章使用technical-blog-writing skills来编写,比起以往的编辑方式,现在真是越来越高效了

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐