O MaxCompute permite que mecanismos de terceiros, como Spark no EMR, StarRocks, Presto, PAI e Hologres, usem um SDK para chamar a Storage API e acessar dados do MaxCompute diretamente. Este tópico fornece exemplos de código para acessar o MaxCompute com o Java SDK.
Visão geral
A tabela a seguir lista as principais interfaces para acessar o MaxCompute com o Java SDK.
|
Interface principal |
Descrição |
|
Cria uma sessão de leitura de tabela do MaxCompute. |
|
|
Representa uma sessão para leitura de dados de uma tabela do MaxCompute. |
|
|
Lê uma partição de dados incluída em uma sessão de leitura de dados. |
Se você usar Maven, pesquise por odps-sdk-table-api no repositório Maven para obter diferentes versões do Java SDK. A configuração relacionada é a seguinte:
<dependency>
<groupId>com.aliyun.odps</groupId>
<artifactId>odps-sdk-table-api</artifactId>
<version>0.48.8-public</version>
</dependency>
O MaxCompute fornece APIs relacionadas ao open storage. Para mais informações, consulte odps-sdk-table-api.
TableReadSessionBuilder
A interface TableReadSessionBuilder cria uma sessão de leitura de tabela do MaxCompute. Os principais métodos são definidos conforme abaixo. Para mais detalhes, consulte Java-sdk-doc.
Definição da interface
public class TableReadSessionBuilder {
public TableReadSessionBuilder table(Table table);
public TableReadSessionBuilder identifier(TableIdentifier identifier);
public TableReadSessionBuilder requiredDataColumns(List<String> requiredDataColumns);
public TableReadSessionBuilder requiredPartitionColumns(List<String> requiredPartitionColumns);
public TableReadSessionBuilder requiredPartitions(List<PartitionSpec> requiredPartitions);
public TableReadSessionBuilder requiredBucketIds(List<Integer> requiredBucketIds);
public TableReadSessionBuilder withSplitOptions(SplitOptions splitOptions);
public TableReadSessionBuilder withArrowOptions(ArrowOptions arrowOptions);
public TableReadSessionBuilder withFilterPredicate(Predicate filterPredicate);
public TableReadSessionBuilder withSettings(EnvironmentSettings settings);
public TableReadSessionBuilder withSessionId(String sessionId);
public TableBatchReadSession buildBatchReadSession();
}
Descrição dos métodos
Nome do método | Descrição |
| Define o parâmetro Table de entrada como a tabela de destino da sessão atual. |
| Define o parâmetro TableIdentifier de entrada como a tabela de destino da sessão atual. |
| Lê dados de campos específicos. A ordem dos campos nos dados retornados corresponde à ordem especificada no parâmetro Nota Se o parâmetro |
| Lê dados de colunas específicas em partições determinadas de uma tabela. Use este método para poda de partições. Nota Caso o parâmetro |
| Lê dados de partições específicas de uma tabela. Recomendado para cenários de poda de partições. Nota Quando o parâmetro |
| Lê dados de buckets específicos. Este método aplica-se apenas a tabelas clusterizadas e serve para poda de buckets. Nota Se o parâmetro |
| Divide os dados da tabela. Para mais informações, consulte SplitOptions. |
| Especifica opções de dados Arrow. Consulte ArrowOptions para detalhes. |
| Define opções de pushdown de predicado. Veja Predicate para mais informações. |
| Especifica o contexto de ambiente. Consulte EnvironmentSettings para saber mais. |
| Define o ID da sessão para recarregar uma sessão existente. |
| Cria ou recupera uma sessão de leitura de tabela.
Nota A criação de uma sessão gera alta sobrecarga e pode levar muito tempo quando o número de arquivos é grande. |
TableBatchReadSession
A interface TableBatchReadSession representa uma sessão para leitura de dados de uma tabela do MaxCompute. Os principais métodos estão definidos a seguir.
Definição da interface
public interface TableBatchReadSession {
String getId();
TableIdentifier getTableIdentifier();
SessionStatus getStatus();
DataSchema readSchema();
InputSplitAssigner getInputSplitAssigner() throws IOException;
SplitReader<ArrayRecord> createRecordReader(InputSplit split, ReaderOptions options) throws IOException;
SplitReader<VectorSchemaRoot> createArrowReader(InputSplit split, ReaderOptions options) throws IOException;
}
Descrição dos métodos
Nome do método | Descrição |
| Obtém o ID da sessão. O tempo limite padrão da sessão é de 24 horas. |
| Retorna o nome da tabela da sessão atual. |
| Obtém o status da sessão. Valores válidos:
|
| Recupera o esquema da tabela para a sessão atual. Para mais informações, veja DataSchema. |
| Obtém o InputSplitAssigner da sessão atual. O InputSplitAssigner define métodos para atribuir instâncias de InputSplit na sessão de leitura atual. Cada InputSplit representa uma partição de dados que um único SplitReader pode processar. Consulte InputSplitAssigner para detalhes. |
| Constrói um objeto |
| Constrói um objeto |
SplitReader
A interface SplitReader lê dados das tabelas.
Definição da interface
public interface SplitReader<T> {
boolean hasNext() throws IOException;
T get();
Metrics currentMetricsValues();
void close() throws IOException;
}
Descrição dos métodos
|
Nome do método |
Descrição |
|
|
Verifica se há mais itens de dados disponíveis. Retorna true se outro item puder ser lido; caso contrário, retorna false. |
|
|
Obtém o item de dados atual. Chame |
|
|
Recupera métricas relacionadas ao SplitReader. |
|
|
Fecha a conexão após o término da leitura. |
Exemplos
-
Configure o ambiente para conexão com o MaxCompute..
// AccessKey ID and AccessKey secret of an Alibaba Cloud account or a RAM user // An Alibaba Cloud account AccessKey grants full API access and poses high security risks. We strongly recommend creating and using a RAM user for API access or routine O&M. Log on to the RAM console to create a RAM user. // This example stores the AccessKey and AccessKey secret in environment variables. They can also be stored in a configuration file as needed. // Never store AccessKey and AccessKey secret in code because of the risk of key leakage. private static String accessId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"); private static String accessKey = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"); // Quota name for accessing MaxCompute String quotaName = "<quotaName>"; // MaxCompute project name String project = "<project>"; // Create an Odps object to connect to MaxCompute Account account = new AliyunAccount(accessId, accessKey); Odps odps = new Odps(account); odps.setDefaultProject(project); // Endpoint for MaxCompute. Only Alibaba Cloud VPC networks are supported. odps.setEndpoint(endpoint); Credentials credentials = Credentials.newBuilder().withAccount(odps.getAccount()).withAppAccount(odps.getAppAccount()).build(); EnvironmentSettings settings = EnvironmentSettings.newBuilder().withCredentials(credentials).withServiceEndpoint(odps.getEndpoint()).withQuotaName(quotaName).build(); -
Autorização.
Por padrão, nenhuma conta (incluindo contas Alibaba Cloud) ou função possui permissões para especificar cotas no nível de job. Conceda as permissões necessárias. Para mais informações, consulte Visão geral do Open Storage.
-
Leia os dados da tabela.
-
Crie uma sessão de leitura de dados para acessar dados do MaxCompute.
// Table name in the MaxCompute project String tableName = "<table.name>"; // Create a table data read session TableReadSessionBuilder scanBuilder = new TableReadSessionBuilder(); TableBatchReadSession scan = scanBuilder.identifier(TableIdentifier.of(project, tableName)).withSettings(settings) .withSplitOptions(SplitOptions.newBuilder() .SplitByByteSize(256 * 1024L * 1024L) .withCrossPartition(false).build()) .requiredDataColumns(Arrays.asList("timestamp")) .requiredPartitionColumns(Arrays.asList("pt1")) .buildBatchReadSession();NotaSe o volume de dados for grande ou a latência da rede for alta ou instável, a criação de uma sessão de leitura de dados pode demorar excessivamente e mudar automaticamente para um processo assíncrono.
-
Percorra os dados do MaxCompute em cada partição e use um leitor Arrow para ler e exibir os dados de cada partição.
// Traverse all input partitions, use an Arrow reader to read each batch of data from every partition, and output the content of each batch InputSplitAssigner assigner = scan.getInputSplitAssigner(); for (InputSplit split : assigner.getAllSplits()) { SplitReader<VectorSchemaRoot> reader = scan.createArrowReader(split, ReaderOptions.newBuilder() .withSettings(settings) .withCompressionCodec(CompressionCodec.ZSTD) .withReuseBatch(true) .build()); int rowCount = 0; List<VectorSchemaRoot> batchList = new ArrayList<>(); while (reader.hasNext()) { VectorSchemaRoot data = reader.get(); rowCount += data.getRowCount(); System.out.println(data.contentToTSVString()); } reader.close(); }
-
Objetos e interfaces relacionados
SplitOptions
SplitOptions
-
Definição de parâmetros
Os parâmetros do objeto SplitOptions são definidos da seguinte forma:
public class SplitOptions { public static SplitOptions.Builder newBuilder() { return new Builder(); } public static class Builder { public SplitOptions.Builder SplitByByteSize(long splitByteSize); public SplitOptions.Builder SplitByRowOffset(); public SplitOptions.Builder withCrossPartition(boolean crossPartition); public SplitOptions.Builder withMaxFileNum(int splitMaxFileNum); public SplitOptions build(); } } -
Descrição dos parâmetros
-
SplitByByteSize(long splitByteSize)Divide os dados com base no parâmetro splitByteSize especificado. O tamanho de cada partição de dados retornada pelo servidor não excede splitByteSize (em bytes).
O tamanho personalizado da divisão deve ser de pelo menos 10 × 1024 × 1024 (10 MB).
Se
SplitByByteSize(long splitByteSize)não for usado para personalizar o tamanho da divisão, o sistema utiliza o valor padrão de 256 × 1024 × 1024 (256 MB).
-
SplitByRowOffset()Divide os dados por linha, permitindo que o cliente leia dados a partir de um índice específico.
-
withCrossPartition(boolean crossPartition)Especifica se permite que um único shard de dados inclua múltiplas partições de dados. O parâmetro crossPartition aceita os seguintes valores:
true (padrão): permite que um único shard de dados contenha múltiplas partições de dados.
false: Não permite.
-
withMaxFileNum(int splitMaxFileNum)Quando uma tabela possui muitos arquivos, especifique o número máximo de arquivos físicos em uma única partição de dados para gerar mais partições de dados.
Por padrão, não há limite para o número de arquivos físicos em uma única partição de dados.
build(): Cria um objeto SplitOptions.
-
-
Exemplos
// 1. Split data by size, set SplitSize to 256 MB SplitOptions splitOptionsByteSize = SplitOptions.newBuilder().SplitByByteSize(256 * 1024L * 1024L).build() // 2. Split data by RowOffset SplitOptions splitOptionsCount = SplitOptions.newBuilder().SplitByRowOffset().build() // 3. Set the maximum number of files in a single split to 1 SplitOptions splitOptionsCount = SplitOptions.newBuilder().SplitByRowOffset().withMaxFileNum(1).build()
ArrowOptions
ArrowOptions
-
Definição de parâmetros
Os parâmetros do objeto ArrowOptions são definidos da seguinte forma:
public class ArrowOptions { public static Builder newBuilder() { return new Builder(); } public static class Builder { public Builder withTimestampUnit(TimestampUnit unit); public Builder withDatetimeUnit(TimestampUnit unit); public ArrowOptions build(); } public enum TimestampUnit { SECOND, MILLI, MICRO, NANO; } } -
Descrição dos parâmetros
-
TimestampUnitEspecifica a unidade para os tipos de dados Timestamp e Datetime. Valores válidos:
SECOND: segundos (s)
MILLI: milissegundos (ms)
MICRO: microssegundos (μs)
NANO: nanossegundos (ns)
-
withTimestampUnit(TimestampUnit unit)Define a unidade para o tipo de dado Timestamp. Padrão: NANO.
-
withDatetimeUnit(TimestampUnit unit)Define a unidade para o tipo de dado Datetime. Padrão: MILLI.
-
Exemplos
ArrowOptions options = ArrowOptions.newBuilder() .withDatetimeUnit(ArrowOptions.TimestampUnit.MILLI) .withTimestampUnit(ArrowOptions.TimestampUnit.NANO) .build() -
Definição de parâmetros
Os parâmetros do objeto Predicate são definidos da seguinte forma:
// 1. Binary operations public class BinaryPredicate extends Predicate { public enum Operator { /** * Binary operation operators */ EQUALS("="), NOT_EQUALS("!="), GREATER_THAN(">"), LESS_THAN("<"), GREATER_THAN_OR_EQUAL(">="), LESS_THAN_OR_EQUAL("<="); } public BinaryPredicate(Operator operator, Serializable leftOperand, Serializable rightOperand); public static BinaryPredicate equals(Serializable leftOperand, Serializable rightOperand); public static BinaryPredicate notEquals(Serializable leftOperand, Serializable rightOperand); public static BinaryPredicate greaterThan(Serializable leftOperand, Serializable rightOperand); public static BinaryPredicate lessThan(Serializable leftOperand, Serializable rightOperand); public static BinaryPredicate greaterThanOrEqual(Serializable leftOperand, Serializable rightOperand); public static BinaryPredicate lessThanOrEqual(Serializable leftOperand, Serializable rightOperand); } // 2. Unary operations public class UnaryPredicate extends Predicate { public enum Operator { /** * Unary operation operators */ IS_NULL("is null"), NOT_NULL("is not null"); } public static UnaryPredicate isNull(Serializable operand); public static UnaryPredicate notNull(Serializable operand); } ### 3. IN and NOT IN public class InPredicate extends Predicate { public enum Operator { /** * IN and NOT IN operators for set membership check */ IN("in"), NOT_IN("not in"); } public InPredicate(Operator operator, Serializable operand, List<Serializable> set); public static InPredicate in(Serializable operand, List<Serializable> set); public static InPredicate notIn(Serializable operand, List<Serializable> set); } // 4. Column names public class Attribute extends Predicate { public Attribute(Object value); public static Attribute of(Object value); } // 5. Constants public class Constant extends Predicate { public Constant(Object value); public static Constant of(Object value); } // 6. Compound operations public class CompoundPredicate extends Predicate { public enum Operator { /** * Compound predicate operators */ AND("and"), OR("or"), NOT("not"); } public CompoundPredicate(Operator logicalOperator, List<Predicate> predicates); public static CompoundPredicate and(Predicate... predicates); public static CompoundPredicate or(Predicate... predicates); public static CompoundPredicate not(Predicate predicates); public void addPredicate(Predicate predicate); } // 7. Raw predicates (RawPredicate) // If existing methods do not meet requirements, assemble predicates based on SQL syntax public class RawPredicate extends Predicate { public RawPredicate(String rawExpr); public static RawPredicate of(String rawExpr); } -
Exemplos
// 1. c1 > 20000 and c2 < 100000 BinaryPredicate c1 = new BinaryPredicate(BinaryPredicate.Operator.GREATER_THAN, Attribute.of("c1"), Constant.of(20000)); BinaryPredicate c2 = new BinaryPredicate(BinaryPredicate.Operator.LESS_THAN, Attribute.of("c2"), Constant.of(100000)); CompoundPredicate predicate = new CompoundPredicate(CompoundPredicate.Operator.AND, ImmutableList.of(c1, c2)); // 2. c1 is not null Predicate predicate = new UnaryPredicate(UnaryPredicate.Operator.NOT_NULL, Attribute.of("c1")); // 3. c1 in (1, 10001) Predicate predicate = new InPredicate(InPredicate.Operator.IN, Attribute.of("c1"), ImmutableList.of(Constant.of(1), Constant.of(10001))); // 4. Use RawPredicate to assemble predicates (supports all types) Predicate predicate = RawPredicate.of("c1 > 20000 and c2 < 100000"); -
Definição de parâmetros
A interface EnvironmentSettings é definida da seguinte forma:
public class EnvironmentSettings { public static Builder newBuilder() { return new Builder(); } public static class Builder { public Builder withDefaultProject(String projectName); public Builder withDefaultSchema(String schema); public Builder withServiceEndpoint(String endPoint); public Builder withTunnelEndpoint(String tunnelEndPoint); public Builder withQuotaName(String quotaName); public Builder withCredentials(Credentials credentials); public Builder withRestOptions(RestOptions restOptions); public EnvironmentSettings build(); } } -
Descrição dos parâmetros
-
withDefaultProject(String projectName)Define o nome do projeto. O parâmetro
projectNameé o nome do projeto MaxCompute.Faça login no console do MaxCompute e alterne a região no canto superior esquerdo.
Escolha para visualizar o nome do projeto MaxCompute.
-
withDefaultSchema(String schema)Define o esquema padrão. O parâmetro schema é o nome do esquema do MaxCompute. Para mais informações sobre esquemas, consulte Operações de esquema.
-
withServiceEndpoint(String endPoint)Define o endpoint do service.Endpoint.
-
withTunnelEndpoint(String tunnelEndPoint)Define o endpoint do túnel.Endpoint.
-
withQuotaName(String quotaName)Especifica o nome da cota a ser utilizada.
O MaxCompute suporta dois tipos de recursos: grupos de recursos exclusivos do Data Transmission Service (assinatura) Obtenha o nome da cota conforme descrito abaixo:
-
Grupo de recursos exclusivo do Data Transmission Service
Faça login no console do MaxCompute e selecione uma região no canto superior esquerdo.
No painel de navegação à esquerda, escolha .
Visualize as cotas disponíveis. Para mais informações, consulte Recursos de computação - Gerenciamento de cotas.
-
Faça login no console do MaxCompute e selecione uma região no canto superior esquerdo.
No painel de navegação à esquerda, escolha .
Na aba Tenant Property, ative a chave Storage API Switch.
-
-
withCredentials(Credentials credentials)Especifica as informações de autenticação. Para mais detalhes, consulte Credentials.
-
-
Definição do objeto
public class Credentials { public static Builder newBuilder() { return new Builder(); } public static class Builder { public Builder withAccount(Account account); public Builder withAppAccount(AppAccount appAccount); public Builder withAppStsAccount(AppStsAccount appStsAccount); public Credentials build(); } } -
Descrição dos parâmetros
-
withAccount(Account account)Especifica o objeto Account do Odps.
-
withAppAccount(AppAccount appAccount)Especifica o objeto appAccount do Odps.
-
withAppStsAccount(AppStsAccount appStsAccount)Especifica o objeto appStsAccount do Odps.
-
withRestOptions(RestOptions restOptions)Especifica a configuração de acesso à rede. RestOptions é definido da seguinte forma:
public class RestOptions implements Serializable { public static Builder newBuilder() { return new RestOptions.Builder(); } public static class Builder { public Builder witUserAgent(String userAgent); public Builder withConnectTimeout(int connectTimeout); public Builder withReadTimeout(int readTimeout); public RestOptions build(); } }witUserAgent(String userAgent): Especifica as informações do userAgent.withConnectTimeout(int connectTimeout): Define o tempo limite de conexão para estabelecer a conexão de rede subjacente. Padrão: 10 segundos.withReadTimeout(int readTimeout): Define o tempo limite de leitura para a conexão de rede subjacente. Padrão: 120 segundos.
-
-
DataSchema é definido da seguinte forma:
public class DataSchema implements Serializable { List<Column> getColumns(); List<String> getPartitionKeys(); List<String> getColumnNames(); List<TypeInfo> getColumnDataTypes(); Optional<Column> getColumn(int columnIndex); Optional<Column> getColumn(String columnName); } -
Descrição dos parâmetros
getColumns(): Obtém informações das colunas da tabela e das partições a serem lidas.getPartitionKeys(): Obtém os nomes das colunas de partição a serem lidas.getColumnNames(): Obtém os nomes das colunas da tabela e das partições a serem lidas.getColumnDataTypes(): Obtém os tipos de dados das colunas da tabela e das partições a serem lidas.getColumn(int columnIndex): Obtém um objeto de coluna por índice. Retorna vazio se o índice estiver fora do intervalo.getColumn(String columnName): Obtém um objeto de coluna por nome. Se o nome da colunacolumnNamenão existir na tabela, retorna vazio.columnName
-
InputSplitAssigner é definido da seguinte forma:
public interface InputSplitAssigner { int getSplitsCount(); long getTotalRowCount(); InputSplit getSplit(int index); InputSplit getSplitByRowOffset(long startIndex, long numRecord); } -
Descrição dos parâmetros
-
getSplitsCount(): Obtém o número de partições de dados na sessão.NotaQuando SplitOptions é SplitByByteSize, este método retorna um valor maior ou igual a 0.
-
getTotalRowCount(): Obtém o número total de linhas de dados na sessão.NotaQuando SplitOptions é SplitByByteSize, este método retorna um valor maior ou igual a 0.
getSplit(int index): Obtém o InputSplit para a partição especificadaIndex. O parâmetroindexvaria de[0,SplitsCount-1].getSplitByRowOffset(long startIndex, long numRecord): Obtém o InputSplit correspondente. Os parâmetros são os seguintes:startIndex: Índice da linha inicial para leitura de dados do InputSplit. Intervalo:[0,RecordCount-1].numRecord: Número de linhas de dados para o InputSplit ler.
-
-
Exemplos
// 1. If SplitOptions is SplitByByteSize TableBatchReadSession scan = ...; InputSplitAssigner assigner = scan.getInputSplitAssigner(); int splitCount = assigner.getSplitsCount(); for (int k = 0; k < splitCount; k++) { InputSplit split = assigner.getSplit(k); ... } // 2. If SplitOptions is SplitByRowOffset TableBatchReadSession scan = ...; InputSplitAssigner assigner = scan.getInputSplitAssigner(); long rowCount = assigner.getTotalRowCount(); long recordsPerSplit = 10000; for (long offset = 0; offset < numRecords; offset += recordsPerSplit) { recordsPerSplit = Math.min(recordsPerSplit, numRecords - offset); InputSplit split = assigner.getSplitByRowOffset(offset, recordsPerSplit); ... } -
ReaderOptions é definido da seguinte forma:
public class ReaderOptions { public static ReaderOptions.Builder newBuilder() { return new Builder(); } public static class Builder { public Builder withMaxBatchRowCount(int maxBatchRowCount); public Builder withMaxBatchRawSize(long batchRawSize); public Builder withCompressionCodec(CompressionCodec codec); public Builder withBufferAllocator(BufferAllocator allocator); public Builder withReuseBatch(boolean reuseBatch); public Builder withSettings(EnvironmentSettings settings); public ReaderOptions build(); } } -
Descrição dos parâmetros
-
withMaxBatchRowCount(int maxBatchRowCount)Especifica o número máximo de linhas por lote retornado pelo servidor. O parâmetro
maxBatchRowCounttem como padrão um máximo de 4096. -
withMaxBatchRawSize(long batchRawSize)Especifica o tamanho máximo bruto em bytes por lote retornado pelo servidor.
-
withCompressionCodec(CompressionCodec codec)Especifica o tipo de compressão de dados. Apenas ZSTD e LZ4_FRAME são suportados.
NotaTransferir grandes quantidades de dados Arrow descompactados diretamente pode aumentar significativamente o tempo de transferência devido aos limites de largura de banda da rede.
Se nenhum tipo de compressão for especificado, os dados não serão compactados por padrão.
-
withBufferAllocator(BufferAllocator allocator)Especifica o alocador de memória para leitura de dados Arrow.
-
withReuseBatch(boolean reuseBatch)Especifica se a memória do ArrowBatch pode ser reutilizada.
reuseBatchValores:true (padrão): A memória do ArrowBatch pode ser reutilizada.
false: A memória do ArrowBatch não pode ser reutilizada.
-
withSettings(EnvironmentSettings settings)Especifica as informações do ambiente de execução.
-
Predicate
Predicate
EnvironmentSettings
EnvironmentSettings
Credentials
Credentials
DataSchema
InputSplitAssigner
ReaderOptions
Referências
Para mais informações sobre o open storage do MaxCompute , consulte Visão geral do Open Storage.
-