Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.
Migrasi konektor Spark Kinesis ke SDK 2.x untuk Amazon EMR 7.0
AWS SDK menyediakan serangkaian API dan pustaka yang kaya untuk berinteraksi dengan layanan komputasi AWS awan, seperti mengelola kredenSIAL, menghubungkan ke layanan S3 dan Kinesis. Konektor Spark Kinesis digunakan untuk mengkonsumsi data dari Kinesis Data Streams, dan data yang diterima diubah dan diproses dalam mesin eksekusi Spark. Saat ini konektor ini dibangun di atas 1.x AWS SDK dan Kinesis-client-library (KCL).
Sebagai bagian dari migrasi AWS SDK 2.x, konektor Spark Kinesis juga diperbarui sesuai untuk dijalankan dengan SDK 2.x. Dalam rilis Amazon EMR 7.0, Spark berisi peningkatan SDK 2.x yang belum tersedia di versi komunitas Apache Spark. Jika Anda menggunakan konektor Spark Kinesis dari rilis yang lebih rendah dari 7.0, Anda harus memigrasikan kode aplikasi untuk dijalankan di SDK 2.x sebelum Anda dapat bermigrasi ke Amazon EMR 7.0.
Panduan migrasi
Bagian ini menjelaskan langkah-langkah untuk memigrasikan aplikasi ke konektor Spark Kinesis yang ditingkatkan. Ini mencakup panduan untuk bermigrasi ke Kinesis Client Library (KCL) 2.x, penyedia AWS kredensia, dan klien AWS layanan di SDK 2.x. AWS Untuk referensi, ini juga termasuk WordCount
Topik
Migrasi KCL dari 1.x ke 2.x
-
Tingkat metrik dan dimensi di
KinesisInputDStreamSaat membuat instance
KinesisInputDStream, Anda dapat mengontrol level metrik dan dimensi untuk streaming. Contoh berikut menunjukkan bagaimana Anda dapat menyesuaikan parameter ini dengan KCL 1.x:import com.amazonaws.services.kinesis.clientlibrary.lib.worker.KinesisClientLibConfiguration import com.amazonaws.services.kinesis.metrics.interfaces.MetricsLevel val kinesisStream = KinesisInputDStream.builder .streamingContext(ssc) .streamName(streamName) .endpointUrl(endpointUrl) .regionName(regionName) .initialPosition(new Latest()) .checkpointAppName(appName) .checkpointInterval(kinesisCheckpointInterval) .storageLevel(StorageLevel.MEMORY_AND_DISK_2) .metricsLevel(MetricsLevel.DETAILED) .metricsEnabledDimensions(KinesisClientLibConfiguration.DEFAULT_METRICS_ENABLED_DIMENSIONS.asScala.toSet) .build()Di KCL 2.x, pengaturan konfigurasi ini memiliki nama paket yang berbeda. Untuk bermigrasi ke 2.x:
-
Ubah pernyataan impor untuk
com.amazonaws.services.kinesis.clientlibrary.lib.worker.KinesisClientLibConfigurationdancom.amazonaws.services.kinesis.metrics.interfaces.MetricsLevelkesoftware.amazon.kinesis.metrics.MetricsLeveldansoftware.amazon.kinesis.metrics.MetricsUtilmasing-masing.// import com.amazonaws.services.kinesis.metrics.interfaces.MetricsLevel import software.amazon.kinesis.metrics.MetricsLevel // import com.amazonaws.services.kinesis.clientlibrary.lib.worker.KinesisClientLibConfiguration import software.amazon.kinesis.metrics.MetricsUtil -
Ganti garis
metricsEnabledDimensionsKinesisClientLibConfiguration.DEFAULT_METRICS_ENABLED_DIMENSIONS.asScala.toSetdenganmetricsEnabledDimensionsSet(MetricsUtil.OPERATION_DIMENSION_NAME, MetricsUtil.SHARD_ID_DIMENSION_NAME)
Berikut ini adalah versi terbaru dari dimensi tingkat metrik dan metrik yang disesuaikan:
KinesisInputDStreamimport software.amazon.kinesis.metrics.MetricsLevel import software.amazon.kinesis.metrics.MetricsUtil val kinesisStream = KinesisInputDStream.builder .streamingContext(ssc) .streamName(streamName) .endpointUrl(endpointUrl) .regionName(regionName) .initialPosition(new Latest()) .checkpointAppName(appName) .checkpointInterval(kinesisCheckpointInterval) .storageLevel(StorageLevel.MEMORY_AND_DISK_2) .metricsLevel(MetricsLevel.DETAILED) .metricsEnabledDimensions(Set(MetricsUtil.OPERATION_DIMENSION_NAME, MetricsUtil.SHARD_ID_DIMENSION_NAME)) .build() -
-
Fungsi penangan pesan di
KinesisInputDStreamSaat membuat instance a
KinesisInputDStream, Anda juga dapat menyediakan “fungsi penangan pesan” yang mengambil Kinesis Record dan mengembalikan objek generik T, jika Anda ingin menggunakan data lain yang disertakan dalam Rekaman seperti kunci partisi.Di KCL 1.x, tanda tangan fungsi pengendali pesan adalah:
Record => T, di mana Record berada.com.amazonaws.services.kinesis.model.RecordDi KCL 2.x, tanda tangan handler diubah menjadi:KinesisClientRecord => T, where KinesisClientRecord is.software.amazon.kinesis.retrieval.KinesisClientRecordBerikut ini adalah contoh penyediaan penangan pesan di KCL 1.x:
import com.amazonaws.services.kinesis.model.Record def addFive(r: Record): Int = JavaUtils.bytesToString(r.getData).toInt + 5 val stream = KinesisInputDStream.builder .streamingContext(ssc) .streamName(streamName) .endpointUrl(endpointUrl) .regionName(regionName) .initialPosition(new Latest()) .checkpointAppName(appName) .checkpointInterval(Seconds(10)) .storageLevel(StorageLevel.MEMORY_ONLY) .buildWithMessageHandler(addFive)Untuk memigrasikan pengendali pesan:
-
Ubah pernyataan impor
com.amazonaws.services.kinesis.model.Recorduntuk menjadisoftware.amazon.kinesis.retrieval.KinesisClientRecord.// import com.amazonaws.services.kinesis.model.Record import software.amazon.kinesis.retrieval.KinesisClientRecord -
Perbarui tanda tangan metode dari penangan pesan.
//def addFive(r: Record): Int = JavaUtils.bytesToString(r.getData).toInt + 5 def addFive = (r: KinesisClientRecord) => JavaUtils.bytesToString(r.data()).toInt + 5
Berikut ini adalah contoh yang diperbarui untuk menyediakan pengendali pesan di KCL 2.x:
import software.amazon.kinesis.retrieval.KinesisClientRecord def addFive = (r: KinesisClientRecord) => JavaUtils.bytesToString(r.data()).toInt + 5 val stream = KinesisInputDStream.builder .streamingContext(ssc) .streamName(streamName) .endpointUrl(endpointUrl) .regionName(regionName) .initialPosition(new Latest()) .checkpointAppName(appName) .checkpointInterval(Seconds(10)) .storageLevel(StorageLevel.MEMORY_ONLY) .buildWithMessageHandler(addFive)Untuk informasi selengkapnya tentang migrasi dari KCL 1.x ke 2.x, lihat Migrasi Konsumen dari KCL 1.x ke KCL 2.x.
-
Migrating AWS penyedia kredenSIAL dari AWS SDK 1.x hingga 2.x
Penyedia kredenSIAL digunakan untuk mendapatkan AWS kredenSIAL untuk interaksi dengan AWS. Ada beberapa perubahan antarmuka dan kelas yang terkait dengan penyedia kredenSIAL di SDK 2.x, yang dapat ditemukan di siniorg.apache.spark.streaming.kinesis.SparkAWSCredentials) dan kelas implementasi yang mengembalikan versi 1.x dari penyedia AWS kredensia. Penyedia kredenSIAL ini diperlukan saat menginisialisasi klien Kinesis. Misalnya, jika Anda menggunakan metode SparkAWSCredentials.provider dalam aplikasi, Anda perlu memperbarui kode untuk menggunakan penyedia AWS kredensia versi 2.x.
Berikut ini adalah contoh penggunaan penyedia kredensia di AWS SDK 1.x:
import org.apache.spark.streaming.kinesis.SparkAWSCredentials import com.amazonaws.auth.AWSCredentialsProvider val basicSparkCredentials = SparkAWSCredentials.builder .basicCredentials("accessKey", "secretKey") .build() val credentialProvider = basicSparkCredentials.provider assert(credentialProvider.isInstanceOf[AWSCredentialsProvider], "Type should be AWSCredentialsProvider")
Untuk bermigrasi ke SDK 2.x:
-
Ubah pernyataan impor untuk
com.amazonaws.auth.AWSCredentialsProvidermenjadisoftware.amazon.awssdk.auth.credentials.AwsCredentialsProvider//import com.amazonaws.auth.AWSCredentialsProvider import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider -
Perbarui kode yang tersisa yang menggunakan kelas ini.
import org.apache.spark.streaming.kinesis.SparkAWSCredentials import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider val basicSparkCredentials = SparkAWSCredentials.builder .basicCredentials("accessKey", "secretKey") .build() val credentialProvider = basicSparkCredentials.provider assert (credentialProvider.isInstanceOf[AwsCredentialsProvider], "Type should be AwsCredentialsProvider")
Migrating AWS Pelayanan Pelanggan dari AWS SDK 1.x hingga 2.x
AWS klien layanan memiliki nama paket yang berbeda di 2.x (yaitusoftware.amazon.awssdk). sedangkan SDK 1.x menggunakan. com.amazonaws Untuk informasi lebih lanjut tentang perubahan klien, lihat di sini. Jika Anda menggunakan klien layanan ini dalam kode, Anda perlu memigrasikan klien yang sesuai.
Berikut ini adalah contoh membuat klien di SDK 1.x:
import com.amazonaws.services.dynamodbv2.AmazonDynamoDBClient import com.amazonaws.services.dynamodbv2.document.DynamoDB AmazonDynamoDB ddbClient = AmazonDynamoDBClientBuilder.defaultClient(); AmazonDynamoDBClient ddbClient = new AmazonDynamoDBClient();
Untuk bermigrasi ke 2.x:
-
Ubah pernyataan impor untuk klien layanan. Ambil klien DynamoDB sebagai contoh. Anda perlu mengubah
com.amazonaws.services.dynamodbv2.AmazonDynamoDBClientataucom.amazonaws.services.dynamodbv2.document.DynamoDBmengubahsoftware.amazon.awssdk.services.dynamodb.DynamoDbClient.// import com.amazonaws.services.dynamodbv2.AmazonDynamoDBClient // import com.amazonaws.services.dynamodbv2.document.DynamoDB import software.amazon.awssdk.services.dynamodb.DynamoDbClient -
Perbarui kode yang menginisialisasi klien
// AmazonDynamoDB ddbClient = AmazonDynamoDBClientBuilder.defaultClient(); // AmazonDynamoDBClient ddbClient = new AmazonDynamoDBClient(); DynamoDbClient ddbClient = DynamoDbClient.create(); DynamoDbClient ddbClient = DynamoDbClient.builder().build();Untuk informasi selengkapnya tentang migrasi AWS SDK dari 1.x ke 2.x, lihat Perbedaan antara AWS SDK untuk Java 1.x dan 2.x
Contoh kode untuk aplikasi streaming
import java.net.URI import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider import software.amazon.awssdk.http.apache.ApacheHttpClient import software.amazon.awssdk.services.kinesis.KinesisClient import software.amazon.awssdk.services.kinesis.model.DescribeStreamRequest import software.amazon.awssdk.regions.Region import software.amazon.kinesis.metrics.{MetricsLevel, MetricsUtil} import org.apache.spark.SparkConf import org.apache.spark.storage.StorageLevel import org.apache.spark.streaming.{Milliseconds, StreamingContext} import org.apache.spark.streaming.dstream.DStream.toPairDStreamFunctions import org.apache.spark.streaming.kinesis.KinesisInitialPositions.Latest import org.apache.spark.streaming.kinesis.KinesisInputDStream object KinesisWordCountASLSDKV2 { def main(args: Array[String]): Unit = { val appName = "demo-app" val streamName = "demo-kinesis-test" val endpointUrl = "https://kinesis.us-west-2.amazonaws.com" val regionName = "us-west-2" // Determine the number of shards from the stream using the low-level Kinesis Client // from the AWS Java SDK. val credentialsProvider = DefaultCredentialsProvider.create require(credentialsProvider.resolveCredentials() != null, "No AWS credentials found. Please specify credentials using one of the methods specified " + "in https://docs.aws.amazon.com/sdk-for-java/latest/developer-guide/credentials.html") val kinesisClient = KinesisClient.builder() .credentialsProvider(credentialsProvider) .region(Region.US_WEST_2) .endpointOverride(URI.create(endpointUrl)) .httpClientBuilder(ApacheHttpClient.builder()) .build() val describeStreamRequest = DescribeStreamRequest.builder() .streamName(streamName) .build() val numShards = kinesisClient.describeStream(describeStreamRequest) .streamDescription .shards .size // In this example, we are going to create 1 Kinesis Receiver/input DStream for each shard. // This is not a necessity; if there are less receivers/DStreams than the number of shards, // then the shards will be automatically distributed among the receivers and each receiver // will receive data from multiple shards. val numStreams = numShards // Spark Streaming batch interval val batchInterval = Milliseconds(2000) // Kinesis checkpoint interval is the interval at which the DynamoDB is updated with information // on sequence number of records that have been received. Same as batchInterval for this // example. val kinesisCheckpointInterval = batchInterval // Setup the SparkConfig and StreamingContext val sparkConfig = new SparkConf().setAppName("KinesisWordCountASLSDKV2") val ssc = new StreamingContext(sparkConfig, batchInterval) // Create the Kinesis DStreams val kinesisStreams = (0 until numStreams).map { i => KinesisInputDStream.builder .streamingContext(ssc) .streamName(streamName) .endpointUrl(endpointUrl) .regionName(regionName) .initialPosition(new Latest()) .checkpointAppName(appName) .checkpointInterval(kinesisCheckpointInterval) .storageLevel(StorageLevel.MEMORY_AND_DISK_2) .metricsLevel(MetricsLevel.DETAILED) .metricsEnabledDimensions(Set(MetricsUtil.OPERATION_DIMENSION_NAME, MetricsUtil.SHARD_ID_DIMENSION_NAME)) .build() } // Union all the streams val unionStreams = ssc.union(kinesisStreams) // Convert each line of Array[Byte] to String, and split into words val words = unionStreams.flatMap(byteArray => new String(byteArray).split(" ")) // Map each word to a (word, 1) tuple so we can reduce by key to count the words val wordCounts = words.map(word => (word, 1)).reduceByKey(_ + _) // Print the first 10 wordCounts wordCounts.print() // Start the streaming context and await termination ssc.start() ssc.awaitTermination() } }
Pertimbangan saat menggunakan konektor Spark Kinesis yang ditingkatkan
-
Jika aplikasi Anda menggunakan versi JDK yang lebih rendah dari 11, Anda mungkin mengalami pengecualian seperti
java.lang.NoClassDefFoundError: javax/xml/bind/DatatypeConverter.Kinesis-producer-libraryIni terjadi karena EMR 7.0 hadir dengan JDK 17 secara default dan modul J2EE telah dihapus dari pustaka standar sejak Java 11+. Ini dapat diperbaiki dengan menambahkan ketergantungan berikut di file pom. Ganti versi pustaka dengan yang Anda inginkan.<dependency> <groupId>javax.xml.bind</groupId> <artifactId>jaxb-api</artifactId> <version>${jaxb-api.version}</version> </dependency> -
Toples konektor Spark Kinesis dapat ditemukan di bawah jalur ini setelah cluster EMR dibuat:
/usr/lib/spark/connector/lib/