Aléatoire désagrégé

Le service de mélange désagrégé permet à la majeure partie des données de mélange de résider dans un emplacement de stockage distinct (montage de volume partagé distinct ou librairie) à partir des noeuds de traitement, ce qui permet une utilisation efficace des ressources. Le plug-in Disaggregated shuffle utilise l'API Spark shuffle manager qui remplace le gestionnaire Shuffle existant. Il enregistre un nouveau composant ShuffleManager et ShuffleDataIO qui permet au plugin de surcharger les méthodes d'E/S Spark shuffle.

Le plug-in Disaggregated shuffle est utilisé par les travaux Spark qui sont à court de stockage temporaire local ou qui connaissent des échecs de pod. Il est également utilisé lorsque les ressources de stockage locales sont insuffisantes ou comme alternative à emptyDir pour améliorer les performances. Le plug-in Disaggregated shuffle prend en charge les stockages suivants:

  • Volumes partagés Cloud Pak for Data
  • Tout stockage compatible Spark tel que cos://bucket, s3a://bucket, hdfs: // ou file:///path
  1. Sauvegardez l'exemple de fichier Python suivant.

    from pyspark.sql import SparkSession
    import time
    
    def init_spark():
    
        spark = SparkSession.builder.appName("auto-scale-test").getOrCreate()
    
        sc = spark.sparkContext
    
        return spark, sc
    
    def main():
    
        spark, sc = init_spark()
    
        partitions = [10, 5, 8, 4, 9, 4, 7]
    
        parallelism = sc.defaultParallelism
    
        for num_partitions in partitions:
    
            print(
    
                f"Running experiment with {num_partitions} partitions leveraging {parallelism} cores."
    
            )
    
            data = range(1, 20000000)
    
            v0 = sc.parallelize(data, parallelism)
    
            gr = v0.groupBy(lambda x: x % num_partitions, numPartitions=num_partitions).map(
    
                lambda x: (x[0], len(x[1]))
    
            )
    
            dd = gr.collect()
    
            print(f"Buckets {dd}")
    
            time.sleep(5)
    
    if __name__ == "__main__":
    
        main()
    
  2. Téléchargez le fichier Python dans le compartiment Cloud Object Storage .

  3. Pour activer la fonction de mélange sur un volume partagé dans Analytics Engine, vous devez monter le volume partagé et spécifier le chemin dans la fonction de mélange désagrégé. Exécutez la commande curl suivante pour soumettre l'application Spark.

curl -k -X POST https://<cluster-url>/v4/analytics_engines/<app-id>/spark_applications \
    -H "Authorization: ZenApiKey <api-key>" \
    -d '{  "application_details": {\
    "conf": {\
           "spark.hadoop.fs.s3a.endpoint": "s3.us-south.cloud-object-storage.appdomain.cloud",\
           "spark.hadoop.fs.s3a.access.key": "<access-key>",\
           "spark.hadoop.fs.s3a.secret.key": "<secret-key>",\
           "spark.hadoop.fs.s3a.path.style.access": "true",\
           "spark.hadoop.fs.s3a.fast.upload": "true","spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",\
           "spark.shuffle.manager": "org.apache.spark.shuffle.sort.DisaggregatedShuffleManager",\
           "spark.shuffle.sort.io.plugin.class": "org.apache.spark.shuffle.DisaggregatedShuffleDataIO",\
           "spark.shuffle.disaggregated.rootDir": "file:///<path-name>/",\
           "spark.shuffle.disaggregated.folderPrefixes": "1"    },\
           "application": "s3a://<cos-bucket-name>/<app-name>.py"  },
            "volumes": [{    "name": "cpd-instance::<volume-name>",
            "mount_path": "file:///<path-name>"  }]}'

Le plug-in nettoie automatiquement les fichiers shuffle créés. Cependant, les fichiers de remaniement créés lorsqu'un pilote Spark est arrêté ou que des pannes ne sont pas removed.It de nettoyer régulièrement le répertoire de remaniement. Si les fichiers shuffle ont une sauvegarde de librairie, la règle de cycle de vie peut être implémentée sur le compartiment.

Optimisations des performances

  • Si les fichiers shuffle sont stockés sur un système de fichiers distribué, il est recommandé de définir spark.shuffle.disaggregated.folderPrefixes sur 1.

  • Pour les librairies telles que IBM COS ou AWS S3, il est recommandé de définir spark.shuffle.disaggregated.folderPrefixes sur 10. Pour de meilleures performances sur les librairies, le plug-in shuffle doit être configuré avec le compartiment de stockage et sans aucun chemin supplémentaire ajouté. Par exemple, si le compartiment est appelé "spark-shuffle-data, spark.shuffle.disaggregated.rootDir doit être défini sur"s3a://spark-shuffle-data/" car l'option de préfixe de dossier permet au plug-in Disaggregé de distribuer des lectures et des écritures sur plusieurs préfixes.

  • Le plug-in disaggreagated shuffle utilise des mémoires tampon d'une taille de 8388608 octets (configurées dans spark.shuffle.disaggregated.bufferSize) pour chaque tâche. Les tampons sont écrits par morceaux de 2097152 octets dans le système de fichiers sous-jacent (configuré dans spark.shuffle.disaggregated.bufferChunkSize, qui par défaut est 1/4th de la taille du tampon). Si votre système de fichiers requiert des blocs plus importants, vous pouvez augmenter la taille de la mémoire tampon et la taille du bloc.

  • Les tailles de mémoire tampon pour lire les données peuvent être configurées séparément. Par défaut, le remaniement désagrégé utilise 134217728 octets (128 Mo) par tâche. Vous pouvez modifier cette valeur en configurant spark.shuffle.disaggregated.maxBufferSizeTask.

  • Chaque tâche de remaniement préextrait les fichiers de remaniement en arrière-plan à l'aide d'un pool d'unités d'exécution. L'accès concurrent de l'utilitaire de lecture anticipée est estimé en fonction du temps d'attente d'E-S mesuré. Le nombre maximal d'unités d'exécution peut être configuré en définissant spark.shuffle.disaggregated.maxConcurrencyTask (par défaut: 5).