Tous les produits
Search
Centre de documentation

MaxCompute:Exemples du SDK Open Storage - Java SDK

Dernière mise à jour :Aug 12, 2026

MaxCompute permet à des moteurs tiers tels que Spark sur EMR, StarRocks, Presto, PAI et Hologres d'utiliser un SDK pour appeler l'API Storage et accéder directement aux données MaxCompute. Cette rubrique fournit des exemples de code pour accéder à MaxCompute à l'aide du SDK Java.

Vue d'ensemble

Le tableau suivant répertorie les principales interfaces permettant d'accéder à MaxCompute via le SDK Java.

Interface principale

Description

TableReadSessionBuilder

Crée une session de lecture de table MaxCompute.

TableBatchReadSession

Représente une session de lecture des données d'une table MaxCompute.

SplitReader

Lit une partition de données incluse dans une session de lecture.

Si vous utilisez Maven, recherchez odps-sdk-table-api dans le référentiel Maven pour obtenir les différentes versions du SDK Java. La configuration associée est la suivante.

<dependency>
	<groupId>com.aliyun.odps</groupId>
	<artifactId>odps-sdk-table-api</artifactId>
	<version>0.48.8-public</version>
</dependency>

MaxCompute fournit des API liées à l'open storage. Pour plus d'informations, consultez odps-sdk-table-api.

TableReadSessionBuilder

L'interface TableReadSessionBuilder permet de créer une session de lecture de table MaxCompute. Les principales méthodes sont définies ci-après. Pour plus d'informations, consultez Java-sdk-doc.

Définition de l'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();
}

Description des méthodes

Nom de la méthode

Description

table(Table table)

Définit le paramètre Table d'entrée comme table cible pour la session en cours.

identifier(TableIdentifier identifier)

Définit le paramètre TableIdentifier d'entrée comme table cible pour la session en cours.

requiredDataColumns(List<String> requiredDataColumns)

Lit les données des champs spécifiés. L'ordre des champs dans les données retournées correspond à l'ordre défini dans le paramètre requiredDataColumns. Utilisez cette méthode pour élaguer les champs de données.

Remarque

Si le paramètre requiredDataColumns est vide, toutes les données de partition sont retournées.

requiredPartitionColumns(List<String> requiredPartitionColumns)

Lit les données des colonnes spécifiées dans les partitions indiquées d'une table. Cette méthode sert à l'élagage des partitions.

Remarque

Si le paramètre requiredPartitionColumns est vide, toutes les données de partition sont retournées.

requiredPartitions(List<PartitionSpec> requiredPartitions)

Lit les données des partitions spécifiées d'une table. Recommandée pour l'élagage des partitions.

Remarque

Si le paramètre requiredPartitions est vide, toutes les données de partition sont retournées.

requiredBucketIds(List<Integer> requiredBucketIds)

Lit les données des buckets spécifiés. Cette méthode s'applique uniquement aux tables clusterisées et permet l'élagage des buckets.

Remarque

Si le paramètre requiredBucketIds est vide, toutes les données des buckets sont retournées.

withSplitOptions(SplitOptions splitOptions)

Fractionne les données de la table. Pour plus d'informations, consultez SplitOptions.

withArrowOptions(ArrowOptions arrowOptions)

Spécifie les options de données Arrow. Pour plus d'informations, consultez ArrowOptions.

withFilterPredicate(Predicate filterPredicate)

Définit les options de pushdown de prédicat. Pour plus d'informations, consultez Predicate.

withSettings(EnvironmentSettings settings)

Spécifie le contexte d'environnement. Pour plus d'informations, consultez EnvironmentSettings.

withSessionId(String sessionId)

Indique l'ID de session pour recharger une session existante.

buildBatchReadSession()

Crée ou récupère une session de lecture de table.

  • Si un ID de session est fourni, la session existante est retournée.

  • Sinon, une nouvelle session est créée.

Remarque

La création d'une session entraîne une charge importante et peut prendre beaucoup de temps lorsque le nombre de fichiers est élevé.

TableBatchReadSession

L'interface TableBatchReadSession représente une session de lecture des données d'une table MaxCompute. Les principales méthodes sont définies comme suit.

Définition de l'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;    

}

Description des méthodes

Nom de la méthode

Description

String getId()

Retourne l'ID de session. Le délai d'expiration par défaut de la session est de 24 heures.

getTableIdentifier()

Retourne le nom de la table pour la session en cours.

getStatus()

Retourne l'état de la session. Valeurs valides :

  • INIT : état initial lors de la création d'une session.

  • NORMAL : session créée avec succès.

  • CRITICAL : échec de la création de la session.

  • EXPIRED : session expirée.

readSchema()

Retourne le schéma de table pour la session en cours. Pour plus d'informations, consultez DataSchema.

getInputSplitAssigner()

Retourne l'InputSplitAssigner pour la session en cours. InputSplitAssigner définit les méthodes d'attribution des instances InputSplit dans la session de lecture actuelle. Chaque InputSplit représente une partition de données qu'un seul SplitReader peut traiter. Pour plus d'informations, consultez InputSplitAssigner.

createRecordReader(InputSplit split, ReaderOptions options)

Construit un objet SplitReader<ArrayRecord> . Pour plus d'informations, consultez ReaderOptions.

createArrowReader(InputSplit split, ReaderOptions options)

Construit un objet SplitReader<VectorSchemaRoot>.

SplitReader

L'interface SplitReader permet de lire les données des tables.

Définition de l'interface

public interface SplitReader<T> {

    boolean hasNext() throws IOException;

    T get();

    Metrics currentMetricsValues();

    void close() throws IOException;
}

Description des méthodes

Nom de la méthode

Description

hasNext()

Vérifie si d'autres éléments de données sont disponibles. Retourne true si un autre élément peut être lu, sinon retourne false.

get()

Retourne l'élément de données actuel. Appelez hasNext() avant d'appeler cette méthode pour confirmer qu'un autre élément est disponible.

currentMetricsValues()

Retourne les métriques associées au SplitReader.

close()

Ferme la connexion une fois la lecture terminée.

Exemples

  1. Configurez l'environnement de connexion à 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();
  2. Autorisation.

    Par défaut, aucun compte (y compris les comptes Alibaba Cloud) ni rôle ne dispose des permissions nécessaires pour spécifier des quotas au niveau des tâches. Accordez les permissions requises. Pour plus d'informations, consultez Vue d'ensemble de l'Open Storage.

  3. Lecture des données de la table.

    1. Créez une session de lecture pour accéder aux données 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();
      Remarque

      Si le volume de données est important ou si la latence réseau est élevée ou instable, la création d'une session de lecture peut prendre trop de temps et basculer automatiquement vers un processus asynchrone.

    2. Parcourez les données MaxCompute de chaque partition et utilisez un lecteur Arrow pour lire et afficher les données de chaque partition.

      // 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();
      }

Objets et interfaces associés

SplitOptions

SplitOptions

  • Définition des paramètres

    Les paramètres de l'objet SplitOptions sont définis comme suit :

    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();
        }
    }
  • Description des paramètres

    • SplitByByteSize(long splitByteSize)

      Fractionne les données en fonction du paramètre splitByteSize spécifié. La taille de chaque partition de données retournée par le serveur ne dépasse pas splitByteSize (en octets).

      • La taille de fractionnement personnalisée doit être d'au moins 10 × 1024 × 1024 (10 Mo).

      • Si SplitByByteSize(long splitByteSize) n'est pas utilisé pour personnaliser la taille de fractionnement, le système utilise la valeur par défaut de 256 × 1024 × 1024 (256 Mo).

    • SplitByRowOffset()

      Fractionne les données par ligne, ce qui permet au client de lire les données à partir d'un index spécifié.

    • withCrossPartition(boolean crossPartition)

      Spécifie si un seul fragment de données peut inclure plusieurs partitions. Le paramètre crossPartition accepte les valeurs suivantes :

      • true (par défaut) : autorise un seul fragment de données à contenir plusieurs partitions.

      • false : ne l'autorise pas.

    • withMaxFileNum(int splitMaxFileNum)

      Lorsqu'une table contient de nombreux fichiers, spécifiez le nombre maximal de fichiers physiques dans une seule partition de données afin de générer davantage de partitions.

      Par défaut, il n'y a aucune limite au nombre de fichiers physiques dans une seule partition de données.

    • build() : crée un objet SplitOptions.

  • Exemples

    // 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

  • Définition des paramètres

    Les paramètres de l'objet ArrowOptions sont définis comme suit :

    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;
        }
    }
  • Description des paramètres

    • TimestampUnit

      Spécifie l'unité pour les types de données Timestamp et Datetime. Valeurs valides :

      • SECOND : secondes (s)

      • MILLI : millisecondes (ms)

      • MICRO : microsecondes (μs)

      • NANO : nanosecondes (ns)

    • withTimestampUnit(TimestampUnit unit)

      Définit l'unité pour le type de données Timestamp. Valeur par défaut : NANO.

    • withDatetimeUnit(TimestampUnit unit)

      Définit l'unité pour le type de données Datetime. Valeur par défaut : MILLI.

  • Exemples

    ArrowOptions options = ArrowOptions.newBuilder()
              .withDatetimeUnit(ArrowOptions.TimestampUnit.MILLI)
              .withTimestampUnit(ArrowOptions.TimestampUnit.NANO)
              .build()

Predicate

Predicate

  • Définition des paramètres

    Les paramètres de l'objet Predicate sont définis comme suit :

    // 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);
      }
  • Exemples

    // 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");

EnvironmentSettings

EnvironmentSettings

  • Définition des paramètres

    L'interface EnvironmentSettings est définie comme suit :

    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();
        }
    }
  • Description des paramètres

    • withDefaultProject(String projectName)

      Définit le nom du projet. Le paramètre projectName correspond au nom du projet MaxCompute.

      • Connectez-vous à la console MaxCompute, puis changez de région dans le coin supérieur gauche.

      • Choisissez Manage Configurations > Projects pour afficher le nom du projet MaxCompute.

    • withDefaultSchema(String schema)

      Définit le schéma par défaut. Le paramètre schema correspond au nom du schéma MaxCompute. Pour plus d'informations sur les schémas, consultez Opérations sur les schémas.

    • withServiceEndpoint(String endPoint)

      Définit l'endpoint du service. Endpoint.

    • withTunnelEndpoint(String tunnelEndPoint)

      Définit l'endpoint du tunnel. Endpoint.

    • withQuotaName(String quotaName)

      Spécifie le nom du quota à utiliser.

      MaxCompute prend en charge deux types de ressources : les groupes de ressources exclusifs Data Transmission Service (abonnement) Obtenez le nom du quota comme suit :

      • Groupe de ressources exclusif Data Transmission Service

        Connectez-vous à la console MaxCompute et sélectionnez une région dans le coin supérieur gauche.

        Dans le volet de navigation de gauche, choisissez Manage Configurations > Quotas .

        Consultez les quotas disponibles. Pour plus d'informations, consultez Ressources de calcul - Gestion des quotas.

      • Connectez-vous à la console MaxCompute et sélectionnez une région dans le coin supérieur gauche.

        Dans le volet de navigation de gauche, choisissez Manage Configurations > Tenants .

        Dans l'onglet Tenant Property, activez le commutateur Storage API Switch.

    • withCredentials(Credentials credentials)

      Spécifie les informations d'authentification. Pour plus d'informations, consultez Credentials.

Credentials

Credentials

  • Définition de l'objet

    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();
        }
    
    }
  • Description des paramètres

    • withAccount(Account account)

      Spécifie l'objet Account Odps.

    • withAppAccount(AppAccount appAccount)

      Spécifie l'objet appAccount Odps.

    • withAppStsAccount(AppStsAccount appStsAccount)

      Spécifie l'objet appStsAccount Odps.

    • withRestOptions(RestOptions restOptions)

      Spécifie la configuration d'accès réseau. RestOptions est défini comme suit :

      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) : spécifie les informations userAgent.

      • withConnectTimeout(int connectTimeout) : définit le délai d'attente de connexion pour établir la connexion réseau sous-jacente. Valeur par défaut : 10 secondes.

      • withReadTimeout(int readTimeout) : définit le délai d'attente de lecture pour la connexion réseau sous-jacente. Valeur par défaut : 120 secondes.

DataSchema

  • DataSchema est défini comme suit :

    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);
    
    }
  • Description des paramètres

    • getColumns() : retourne les informations de colonne pour la table et les partitions à lire.

    • getPartitionKeys() : retourne les noms des colonnes de partition à lire.

    • getColumnNames() : retourne les noms de colonne pour la table et les partitions à lire.

    • getColumnDataTypes() : retourne les types de données de colonne pour la table et les partitions à lire.

    • getColumn(int columnIndex) : retourne un objet colonne par index. Retourne vide si l'index est hors limites.

    • getColumn(String columnName) : retourne un objet colonne par nom. Si le nom de colonne columnName n'existe pas dans la table, retourne vide.columnName

InputSplitAssigner

  • InputSplitAssigner est défini comme suit :

    public interface InputSplitAssigner {
    
        int getSplitsCount();
    
        long getTotalRowCount();
    
        InputSplit getSplit(int index);
    
        InputSplit getSplitByRowOffset(long startIndex, long numRecord);
    }
  • Description des paramètres

    • getSplitsCount() : retourne le nombre de partitions de données dans la session.

      Remarque

      Lorsque SplitOptions est défini sur SplitByByteSize, cette méthode retourne une valeur supérieure ou égale à 0.

    • getTotalRowCount() : retourne le nombre total de lignes de données dans la session.

      Remarque

      Lorsque SplitOptions est défini sur SplitByByteSize, cette méthode retourne une valeur supérieure ou égale à 0.

    • getSplit(int index) : retourne l'InputSplit pour la partition spécifiée Index. Le paramètre index varie de [0,SplitsCount-1].

    • getSplitByRowOffset(long startIndex, long numRecord) : retourne l'InputSplit correspondant. Les paramètres sont les suivants :

    • startIndex : index de ligne de départ pour la lecture des données InputSplit. Plage : [0,RecordCount-1].

    • numRecord : nombre de lignes de données à lire par InputSplit.

  • Exemples

    // 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

  • ReaderOptions est défini comme suit :

    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();
        }
    
    }
  • Description des paramètres

    • withMaxBatchRowCount(int maxBatchRowCount)

      Spécifie le nombre maximal de lignes par lot retourné par le serveur. Le paramètre maxBatchRowCount est limité à 4096 par défaut.

    • withMaxBatchRawSize(long batchRawSize)

      Spécifie la taille brute maximale en octets par lot retourné par le serveur.

    • withCompressionCodec(CompressionCodec codec)

      Spécifie le type de compression des données. Seuls ZSTD et LZ4_FRAME sont pris en charge.

      Remarque
      • Le transfert direct de grandes quantités de données Arrow non compressées peut augmenter considérablement le temps de transfert en raison des limites de bande passante réseau.

      • Si aucun type de compression n'est spécifié, les données ne sont pas compressées par défaut.

    • withBufferAllocator(BufferAllocator allocator)

      Spécifie l'allocateur de mémoire pour la lecture des données Arrow.

    • withReuseBatch(boolean reuseBatch)

      Indique si la mémoire ArrowBatch peut être réutilisée. Valeurs de reuseBatch :

      • true (par défaut) : la mémoire ArrowBatch peut être réutilisée.

      • false : la mémoire ArrowBatch ne peut pas être réutilisée.

    • withSettings(EnvironmentSettings settings)

      Spécifie les informations sur l'environnement d'exécution.

Références

Pour plus d'informations sur l'open storage MaxCompute , consultez Vue d'ensemble de l'Open Storage.