O Stream Load é um método síncrono baseado em HTTP para carregar arquivos locais ou fluxos de dados no ApsaraDB for SelectDB. Ao enviar uma requisição HTTP PUT, o resultado retorna imediatamente no corpo da resposta, sem necessidade de polling. Formatos suportados: CSV, JSON, Parquet e ORC.
Verifique sempre o campo Status no corpo da resposta para determinar se a importação foi bem-sucedida. Uma resposta HTTP 200 indica apenas que a requisição foi recebida — o resultado real da importação está no corpo da mensagem.
Pré-requisitos
Antes de começar, certifique-se de ter:
Conectividade de rede entre a máquina que executa o Stream Load e sua instância SelectDB
Credenciais de login (nome de usuário e senha) da instância SelectDB
Configure o acesso à rede:
Se sua máquina não estiver na mesma Virtual Private Cloud (VPC) da instância SelectDB, solicite um endpoint público.
Adicione os endereços IP da sua máquina à lista de permissões da instância.
-
Caso sua máquina possua uma lista de permissões de saída, adicione o intervalo de IPs da instância SelectDB a ela:
IP da VPC: Consulte Como visualizo os endereços IP na VPC à qual minha instância ApsaraDB SelectDB pertence?
IP público: Execute
ping <public-endpoint>para obter o endereço IP.
Observações de uso
Um único job do Stream Load pode gravar de várias centenas de MB até 1 GB de dados. Carregar pequenos volumes de dados com alta frequência degrada o desempenho da instância e pode causar deadlocks nas tabelas. Utilize processamento em lote para reduzir a frequência de carga:
Lote no lado da aplicação: Acumule dados na aplicação antes de enviar uma requisição Stream Load.
Lote no lado do servidor: Use o Group Commit para permitir que o SelectDB agrupe as requisições recebidas no servidor.
Como funciona o Stream Load
Envie uma requisição HTTP PUT com o arquivo de dados anexado.
O SelectDB processa os dados sincronamente e os grava na tabela de destino.
O corpo da resposta contém o resultado da importação, incluindo um campo
Statuse métricas por linha.
Carregar dados com Stream Load
O Stream Load utiliza requisições HTTP PUT. Os exemplos abaixo usam curl, mas qualquer cliente HTTP é compatível.
Sintaxe
curl --location-trusted \
-u <username>:<password> \
-H "expect:100-continue" \
[-H "<header-key>:<header-value>"] \
-T <file-path> \
-XPUT http://<host>:<port>/api/<db_name>/<table_name>/_stream_load
Substitua os placeholders pelos valores reais:
|
Placeholder |
Descrição |
Exemplo |
|
|
Nome de usuário do SelectDB |
|
|
|
Senha do SelectDB |
|
|
|
Endpoint da VPC ou endpoint público da instância |
|
|
|
Porta HTTP. Padrão: |
|
|
|
Nome do banco de dados de destino |
|
|
|
Nome da tabela de destino |
|
|
|
Caminho para o arquivo de dados local |
|
Escolha do endpoint correto:
Mesma VPC: Use o endpoint da VPC.
VPC diferente ou fora da Alibaba Cloud: Use o endpoint público. Ambos os endpoints estão listados na página de detalhes da instância no console SelectDB.
Cabeçalhos da requisição
Defina as opções de importação como cabeçalhos HTTP usando -H "key:value".
|
Cabeçalho |
Padrão |
Descrição |
|
|
Gerado pelo sistema |
ID exclusivo para este job de importação. Use o mesmo label para novas tentativas no mesmo lote de dados e evitar importações duplicadas (semântica At-Most-Once). Labels podem ser reutilizados quando o job correspondente atingir o status |
|
|
|
Formato dos dados. Valores suportados: |
|
|
|
Delimitador de colunas. Suporta delimitadores com múltiplos caracteres. Para caracteres não imprimíveis, use hexadecimal com o prefixo |
|
|
|
Delimitador de linhas. No Windows, use |
|
|
Nenhum |
Formato de compressão. Valores suportados: |
|
|
|
Fração máxima de linhas que podem falhar nas verificações de qualidade de dados. Intervalo: |
|
|
|
Quando |
|
|
Padrão da instância |
Especifica qual cluster de computação processará a importação. Se nenhum cluster padrão estiver definido, o SelectDB escolhe um automaticamente com base nas suas permissões. |
|
|
|
Quando |
|
|
Nenhum |
Condição de filtro SQL. Linhas que não correspondem são excluídas da importação e contadas em |
|
|
Nenhum |
Restringe a importação às partições especificadas. Linhas fora dessas partições são excluídas e contadas em |
|
|
Nenhum |
Mapeia e transforma colunas de origem. Suporta reordenação de colunas e transformações via expressões SQL, usando a mesma sintaxe das expressões SELECT. |
|
|
|
Comportamento de mesclagem de dados. Opções: |
|
|
Nenhum |
Condição SQL para marcar linhas como excluídas. Usado apenas quando |
|
|
Nenhum |
Para tabelas Unique Key com colunas de sequência. Especifica qual coluna (dos dados de origem ou do esquema da tabela) determina a ordem de substituição das linhas. |
|
|
|
Limite de memória para o job de importação, em bytes. Padrão: 2 GiB. |
|
|
|
Tempo limite de importação em segundos. Intervalo: |
|
|
|
Fuso horário para funções relacionadas a tempo durante a importação. Usa nomes de fuso horário IANA. |
|
|
|
Quando |
|
|
Nenhum |
Padrão de extração de campos JSON, por exemplo |
|
|
|
Expressão JSONPath que seleciona um objeto filho como raiz para análise. |
|
|
|
Quando |
|
|
|
Quando |
Exemplo
Importe data.csv para test_table em test_db:
curl --location-trusted \
-u admin:admin_123 \
-T data.csv \
-H "label:123" \
-H "expect:100-continue" \
http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Resposta
O Stream Load retorna o resultado sincronamente. Verifique o campo Status — e não o código de status HTTP — para determinar o sucesso.
{
"TxnId": 17,
"Label": "707717c0-271a-44c5-be0b-4e71bfeacaa5",
"Status": "Success",
"Message": "OK",
"NumberTotalRows": 5,
"NumberLoadedRows": 5,
"NumberFilteredRows": 0,
"NumberUnselectedRows": 0,
"LoadBytes": 28,
"LoadTimeMs": 27,
"BeginTxnTimeMs": 0,
"StreamLoadPutTimeMs": 2,
"ReadDataTimeMs": 0,
"WriteDataTimeMs": 3,
"CommitAndPublishTimeMs": 18
}
|
Campo |
Descrição |
|
|
ID da transação. |
|
|
Label de importação. Personalizado ou gerado pelo sistema. |
|
|
Resultado da importação. |
|
|
Status do job existente para um label duplicado. Presente apenas quando |
|
|
Mensagem de erro, se houver. |
|
|
Total de linhas processadas. |
|
|
Linhas importadas com sucesso. |
|
|
Linhas filtradas devido a problemas de qualidade de dados. |
|
|
Linhas excluídas pela condição |
|
|
Bytes importados. |
|
|
Tempo total de importação, em milissegundos. |
|
|
Tempo para iniciar a transação no frontend (FE), em milissegundos. |
|
|
Tempo para obter o plano de execução do FE, em milissegundos. |
|
|
Tempo gasto lendo os dados de origem, em milissegundos. |
|
|
Tempo gasto gravando dados, em milissegundos. |
|
|
Tempo para confirmar e publicar a transação, em milissegundos. |
|
|
Se houver problemas de qualidade de dados, acesse esta URL para visualizar as linhas com erro específico. |
Inspecionar erros de importação
Caso existam problemas de qualidade de dados, use ErrorURL para baixar as linhas rejeitadas:
curl "<ErrorURL>"
# or save to a file:
wget "<ErrorURL>" -O error_rows.txt
Cancelar e visualizar tarefas de importação
Não é possível cancelar tarefas do Stream Load manualmente após o envio. O sistema cancela automaticamente uma tarefa se ela exceder o tempo limite ou encontrar um erro irrecuperável.
Para visualizar tarefas concluídas do Stream Load, primeiro ative os registros de operação do Stream Load, depois conecte-se à instância usando um cliente MySQL e execute:
SHOW STREAM LOAD;
Importar dados CSV
Exemplo: Importar usando script
Configurar a tabela de destino
-
Crie o banco de dados e a tabela:
CREATE DATABASE test_db; CREATE TABLE test_table ( id int, name varchar(50), age int, address varchar(50), url varchar(500) ) UNIQUE KEY(`id`, `name`) DISTRIBUTED BY HASH(id) BUCKETS 16 PROPERTIES("replication_num" = "1"); -
Na máquina onde executará o Stream Load, crie um arquivo chamado
test.csv:1,yang,32,shanghai,http://example.com 2,wang,22,beijing,http://example.com 3,xiao,23,shenzhen,http://example.com 4,jess,45,hangzhou,http://example.com 5,jack,14,shanghai,http://example.com 6,tomy,25,hangzhou,http://example.com 7,lucy,45,shanghai,http://example.com 8,tengyin,26,shanghai,http://example.com 9,wangli,27,shenzhen,http://example.com 10,xiaohua,37,shanghai,http://example.com
Deduplicar com label e definir tempo limite personalizado
Importe test.csv com um label para evitar importações duplicadas e um tempo limite personalizado de 100 segundos:
curl --location-trusted \
-u admin:admin_123 \
-H "label:123" \
-H "timeout:100" \
-H "expect:100-continue" \
-H "column_separator:," \
-T test.csv \
http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Filtrar linhas por valor de coluna
Importe apenas linhas onde address seja hangzhou, usando mapeamento de colunas para especificar a ordem dos campos:
curl --location-trusted \
-u admin:admin_123 \
-H "label:123" \
-H "columns: id,name,age,address,url" \
-H "where: address='hangzhou'" \
-H "expect:100-continue" \
-H "column_separator:," \
-T test.csv \
http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Permitir tolerância de erro de 20%
Importe permitindo que até 20% das linhas falhem nas verificações de qualidade de dados:
curl --location-trusted \
-u admin:admin_123 \
-H "label:123" \
-H "max_filter_ratio:0.2" \
-H "expect:100-continue" \
-T test.csv \
http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Ativar modo estrito com fuso horário personalizado
Aplique verificação rigorosa de tipos e defina o fuso horário para Africa/Abidjan:
curl --location-trusted \
-u admin:admin_123 \
-H "strict_mode: true" \
-H "timezone: Africa/Abidjan" \
-H "expect:100-continue" \
-T test.csv \
http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Excluir linhas correspondentes no SelectDB
Exclua todas as linhas em test_table que correspondam às linhas em test.csv:
curl --location-trusted \
-u admin:admin_123 \
-H "merge_type: DELETE" \
-H "expect:100-continue" \
-T test.csv \
http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Excluir linhas por condição e importar o restante
Exclua linhas onde address seja hangzhou e importe as linhas restantes:
curl --location-trusted \
-u admin:admin_123 \
-H "expect:100-continue" \
-H "columns: id,name,age,address,url" \
-H "merge_type: MERGE" \
-H "delete: address='hangzhou'" \
-H "column_separator:," \
-T test.csv \
http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/testDb/testTbl/_stream_load
Exemplo: Importar usando código Java
O exemplo a seguir mostra uma implementação completa em Java usando Apache HttpClient. Verifique o campo Status no corpo da resposta para determinar o sucesso — e não o código de status HTTP.
package com.selectdb.x2doris.connector.doris.writer;
import com.alibaba.fastjson2.JSON;
import org.apache.http.HttpHeaders;
import org.apache.http.HttpResponse;
import org.apache.http.HttpStatus;
import org.apache.http.client.HttpClient;
import org.apache.http.client.config.RequestConfig;
import org.apache.http.client.methods.HttpPut;
import org.apache.http.entity.BufferedHttpEntity;
import org.apache.http.entity.StringEntity;
import org.apache.http.impl.client.DefaultRedirectStrategy;
import org.apache.http.impl.client.HttpClientBuilder;
import org.apache.http.impl.client.HttpClients;
import org.apache.http.protocol.RequestContent;
import org.apache.http.util.EntityUtils;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
import java.util.Base64;
import java.util.List;
import java.util.Map;
public class DorisLoadCase {
public static void main(String[] args) throws Exception {
// 1. Configure parameters.
String loadUrl = "http://<Host:Port>/api/<DB>/<TABLE>/_stream_load?";
String userName = "admin";
String password = "****";
// 2. Build the HTTP client. Redirection (isRedirectable) must be enabled.
HttpClientBuilder httpClientBuilder = HttpClients.custom().setRedirectStrategy(new DefaultRedirectStrategy() {
@Override
protected boolean isRedirectable(String method) {
return true;
}
});
httpClientBuilder.addInterceptorLast(new RequestContent(true));
HttpClient httpClient = httpClientBuilder.build();
// 3. Build the PUT request.
HttpPut httpPut = new HttpPut(loadUrl);
// Set authentication and content headers.
String basicAuth = Base64.getEncoder().encodeToString(
String.format("%s:%s", userName, password).getBytes(StandardCharsets.UTF_8));
httpPut.addHeader(HttpHeaders.AUTHORIZATION, "Basic " + basicAuth);
httpPut.addHeader(HttpHeaders.EXPECT, "100-continue");
httpPut.addHeader(HttpHeaders.CONTENT_TYPE, "text/plain; charset=UTF-8");
RequestConfig reqConfig = RequestConfig.custom().setConnectTimeout(30000).build();
httpPut.setConfig(reqConfig);
// 4. Read the CSV file and attach it as the request body.
// Default CSV delimiters: row = \n, column = \t
List<String> lines = Files.readAllLines(Paths.get("your_file.csv"));
String data = String.join("\n", lines);
httpPut.setEntity(new StringEntity(data));
// 5. Send the request and check the result.
HttpResponse httpResponse = httpClient.execute(httpPut);
int httpStatus = httpResponse.getStatusLine().getStatusCode();
String respContent = EntityUtils.toString(
new BufferedHttpEntity(httpResponse.getEntity()), StandardCharsets.UTF_8);
String respMsg = httpResponse.getStatusLine().getReasonPhrase();
if (httpStatus == HttpStatus.SC_OK) {
// HTTP 200 does not guarantee import success.
// Always check the Status field returned by SelectDB.
Map<String, String> respAsMap = JSON.parseObject(respContent, Map.class);
String dorisStatus = respAsMap.get("Status");
List<String> DORIS_SUCCESS_STATUS = Arrays.asList("Success", "Publish Timeout", "200");
if (!DORIS_SUCCESS_STATUS.contains(dorisStatus) || !respMsg.equals("OK")) {
throw new RuntimeException(
"StreamLoad failed, status: " + dorisStatus + ", Response: " + respMsg);
} else {
System.out.println("Import successful.");
}
} else {
throw new IOException(
"StreamLoad HTTP error: " + httpStatus + ", url: " + loadUrl + ", error: " + respMsg);
}
}
}
Importar dados JSON
O formato JSON não-array tem desempenho significativamente superior ao formato array. Use JSON delimitado por linhas (read_json_by_line:true) sempre que possível.
Configurar a tabela de destino
-
Crie o banco de dados e a tabela:
CREATE DATABASE test_db; CREATE TABLE test_table ( id int, name varchar(50), age int ) UNIQUE KEY(`id`) DISTRIBUTED BY HASH(`id`) BUCKETS 16 PROPERTIES("replication_num" = "1");
Importar JSON delimitado por linhas (recomendado)
Crie um arquivo json.data com um objeto JSON por linha:
{"id":1,"name":"Emily","age":25}
{"id":2,"name":"Benjamin","age":35}
{"id":3,"name":"Olivia","age":28}
{"id":4,"name":"Alexander","age":60}
{"id":5,"name":"Ava","age":17}
Importe-o com read_json_by_line:true:
curl --location-trusted \
-u admin:admin_123 \
-H "Expect:100-continue" \
-H "format:json" \
-H "read_json_by_line:true" \
-T json.data \
-XPUT http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Importar um array JSON
Crie um arquivo json_array.data no formato de array JSON:
[
{"userid":1,"username":"Emily","userage":25},
{"userid":2,"username":"Benjamin","userage":35},
{"userid":3,"username":"Olivia","userage":28},
{"userid":4,"username":"Alexander","userage":60},
{"userid":5,"username":"Ava","userage":17}
]
Importe-o com strip_outer_array:true e jsonpaths para mapear os nomes de campos não correspondentes:
curl --location-trusted \
-u admin:admin_123 \
-H "Expect:100-continue" \
-H "format:json" \
-H "jsonpaths:[\"$.userid\", \"$.userage\", \"$.username\"]" \
-H "columns:id,age,name" \
-H "strip_outer_array:true" \
-T json_array.data \
-XPUT http://selectdb-cn-h033cjs****-fe.selectdbfe.pre.rds.aliyuncs.com:8080/api/test_db/test_table/_stream_load
Modo HTTP Stream
O modo HTTP Stream (http_stream) permite especificar parâmetros de importação como uma expressão SQL no cabeçalho da requisição, utilizando o recurso Table Value Function (TVF). O endpoint da API difere do Stream Load padrão:
Stream Load padrão:
http://host:http_port/api/{db}/{table}/_stream_loadModo HTTP Stream:
http://host:http_port/api/_http_stream
Sintaxe
curl --location-trusted \
-u <username>:<password> \
-H "sql: ${load_sql}" \
-T <file_name> \
-XPUT http://host:http_port/api/_http_stream
O cabeçalho load_sql substitui cabeçalhos individuais como column_separator, line_delimiter, where e columns por uma única instrução SQL:
INSERT INTO db.table (col, ...) SELECT stream_col, ... FROM http_stream("property1"="value1");
Exemplo
curl --location-trusted \
-u admin:admin_123 \
-T test.csv \
-H "sql:insert into demo.example_tbl_1(user_id, age, cost) select c1, c4, c7 * 2 from http_stream(\"format\" = \"CSV\", \"column_separator\" = \",\" ) where age >= 30" \
http://host:http_port/api/_http_stream
Para mais informações sobre TVFs, consulte TVF.Formatos de arquivo
Configuração opcional
Opcional: Ativar registros do Stream Load
Por padrão, o cluster de computação não registra operações do Stream Load. Para ativar o registro:
Defina o parâmetro de backend
enable_stream_load_recordcomotrue.Reinicie o cluster de computação.
Ativar este recurso requer a abertura de um ticket de suporte.
Opcional: Aumentar o tamanho máximo de arquivo
O tamanho máximo padrão de arquivo para um único job do Stream Load é 10.240 MB. Para aumentá-lo, ajuste o parâmetro de backend streaming_load_max_mb. Para instruções, consulte Configurar parâmetros.
Opcional: Ajustar o tempo limite padrão
O tempo limite padrão é de 600 segundos. Substitua-o por job usando o cabeçalho timeout. Para alterar o padrão global, defina o parâmetro de frontend (FE) stream_load_default_timeout_second e reinicie a instância.
Alterar parâmetros globais requer a abertura de um ticket de suporte.
FAQ
O que causa o erro "get table cloud commit lock timeout"?
Este erro indica que as gravações ocorrem com muita frequência, causando bloqueio na tabela. Reduza a frequência de gravação e agrupe seus dados em lotes. Um único job do Stream Load deve visar de várias centenas de MB a 1 GB de dados por requisição. Consulte Observações de uso para estratégias de lote.
Como lidar com dados CSV que contêm delimitadores de coluna ou linha?
Especifique novos delimitadores e atualize o arquivo de dados para que os caracteres delimitadores nos seus dados não entrem em conflito com os delimitadores escolhidos.
Dados contêm o delimitador de linha
Se seus dados contiverem o delimitador de linha padrão \n como valor de dado (não como limite de linha), especifique um delimitador de linha diferente.
Exemplo — arquivo original:
Zhang San\n,25,Shaanxi
Li Si\n,30,Beijing
Passos:
Defina um novo delimitador de linha: adicione
-H "line_delimiter:\r\n"à sua requisição.-
Atualize o arquivo para terminar cada linha com o novo delimitador:
Zhang San\n,25,Shaanxi\r\n Li Si\n,30,Beijing\r\n
Dados contêm o delimitador de coluna
Se seus dados contiverem o delimitador de coluna padrão \t (tabulação) como valor de dado, especifique um delimitador de coluna diferente.
Exemplo — arquivo original:
Zhang San\t 25 Shaanxi
Li Si\t 30 Beijing
Passos:
Defina um novo delimitador de coluna: adicione
-H "column_separator:,"à sua requisição.-
Atualize o arquivo para separar colunas com o novo delimitador:
Zhang San\t,25,Shaanxi Li Si\t,30,Beijing
Próximos passos
Group Commit — agrupe requisições Stream Load recebidas no lado do servidor para reduzir a frequência de gravação
Configurar parâmetros — ajuste parâmetros de backend como
streaming_load_max_mbConectar-se a uma instância ApsaraDB for SelectDB usando um cliente MySQL