戰(zhàn):從 Direct Runner 到分布式 Runner 的選型與運(yùn)行指南)
【免費(fèi)下載鏈接】beamApache Beam is a unified programming model for Batch and Streaming data processing.項(xiàng)目地址https://gitcode.com/gh_mirrors/beam18/beam點(diǎn)擊查看免費(fèi)下載本指南基于 learning/tour-of-beam/learning-content/introduction/introduction-concepts/runner-concepts/description.md 整理而成。Apache Beam 提供一套可移植的 API 層用于構(gòu)建復(fù)雜的數(shù)據(jù)并行處理管線Pipeline并允許同一份管線代碼運(yùn)行在多種執(zhí)行引擎即 Runner之上。本文將以該文檔為主線結(jié)合本倉庫中各 Runner 的源碼實(shí)現(xiàn)系統(tǒng)講解 Beam Runner 的核心概念、Direct Runner 的模型校驗(yàn)機(jī)制以及 Google Cloud Dataflow、Apache Flink、Apache Spark、Apache Samza、Apache Nemo、Hazelcast Jet 等 Runner 的適用場景與運(yùn)行方式幫助讀者完成 Runner 選型并掌握在 Java、Python、Go 三種 SDK 下的實(shí)際運(yùn)行命令。1. 核心概念Beam Runner 是什么Apache Beam 提供的可移植 API 層允許開發(fā)者只編寫一次管線代碼就將其執(zhí)行在不同的執(zhí)行引擎Runner上。這一層的核心概念基于 Beam Model即早先所說的 Dataflow Model而每個(gè) Runner 只是對(duì)該模型不同程度的實(shí)現(xiàn)。也就是說管線Pipeline描述數(shù)據(jù)處理的完整邏輯與具體執(zhí)行環(huán)境無關(guān)Runner負(fù)責(zé)把這條管線翻譯并調(diào)度到某個(gè)具體執(zhí)行引擎上例如本地 JVM、Google Cloud 托管服務(wù)、Flink 集群或 Spark 集群不同 Runner 對(duì) Beam Model 的支持程度不同因此同一份代碼在不同 Runner 上的行為可能有細(xì)微差別這也正是 Direct Runner 存在的重要原因。從源碼結(jié)構(gòu)看本倉庫在runners/目錄下按執(zhí)行引擎組織各個(gè) Runner 模塊例如 runners/direct-java、runners/google-cloud-dataflow-java、runners/flink、runners/spark、runners/samza、runners/jet 等每個(gè)模塊內(nèi)部都實(shí)現(xiàn)了 Beam 的PipelineRunner接口。例如 FlinkRunner.java 聲明為public class FlinkRunner extends PipelineRunnerPipelineResultSparkRunner.java 聲明為public final class SparkRunner extends PipelineRunnerSparkPipelineResult可見所有 Runner 都是 PipelineRunner 的實(shí)現(xiàn)這一模型在代碼層面是統(tǒng)一成立的。2. Direct Runner在本地校驗(yàn) Beam 模型的正確性2.1 設(shè)計(jì)目標(biāo)驗(yàn)證語義而非追求性能Direct Runner 在開發(fā)者本機(jī)執(zhí)行管線它的設(shè)計(jì)目標(biāo)不是高效執(zhí)行而是盡可能嚴(yán)格地驗(yàn)證管線是否遵守 Apache Beam 模型防止用戶依賴那些模型并不保證的語義。從 runners/direct-java 模塊的 DirectRunner.java 源碼注釋可以確認(rèn)The DirectRunner is suitable for running a Pipeline on small scale, example, and test data, and should be used for ensuring that processing logic is correct.也就是說Direct Runner 適合小規(guī)模、示例與測試數(shù)據(jù)其核心價(jià)值在于保證處理邏輯正確并為后續(xù)在分布式后端大規(guī)模執(zhí)行掃清隱患。2.2 模型校驗(yàn)的具體檢查項(xiàng)原文檔指出 Direct Runner 會(huì)執(zhí)行以下額外檢查強(qiáng)制元素不可變性immutability任何 Transform 都不允許在任意時(shí)刻修改輸入元素也不允許在輸出元素被產(chǎn)出之后修改它們強(qiáng)制元素可編碼性encodability每個(gè) PCollection 中的全部元素都必須能被該 PCollection 的 Coder 編碼與解碼任意順序處理所有環(huán)節(jié)中的元素都以任意順序被處理用戶函數(shù)可序列化DoFn、CombineFn 等用戶函數(shù)必須能夠被序列化。上述檢查在源碼中有明確對(duì)應(yīng)。在 DirectOptions.java 中isEnforceImmutability()與isEnforceEncodability()兩個(gè)選項(xiàng)默認(rèn)均為true分別控制是否校驗(yàn)元素不被篡改、是否校驗(yàn)元素可被 Coder 編碼解碼。在 DirectRunner.java 中Enforcement枚舉實(shí)現(xiàn)了ENCODABILITY與IMMUTABILITY兩類強(qiáng)制檢查其中不可變性檢查會(huì)通過ImmutabilityCheckingBundleFactory包裹 Bundle 工廠見bundleFactoryFor方法并且只會(huì)作用于包含用戶函數(shù)UDF的變換源碼中限定了READ_TRANSFORM_URN與PAR_DO_TRANSFORM_URN。2.3 為什么本地單元測試如此重要原文檔特別強(qiáng)調(diào)當(dāng)管線運(yùn)行在遠(yuǎn)程集群上時(shí)排查失敗運(yùn)行往往非常困難而在本地用 Direct Runner 對(duì)管線代碼做單元測試則又快又簡單還能使用你熟悉的本地調(diào)試工具。用 Direct Runner 做測試與開發(fā)有助于保證管線在不同 Beam Runner 之間具有魯棒性。2.4 各 SDK 中的使用方式Go SDK在 Go SDK 中默認(rèn) Runner 就是DirectRunner??芍苯影惭b并運(yùn)行官方 wordcount 示例$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount $ wordcount --input PATH_TO_INPUT_FILE --output counts該示例的源碼位于 sdks/go/examples/wordcount/wordcount.go其注釋同樣強(qiáng)調(diào)可以在本地執(zhí)行管線也可以通過選擇其他 Runner 執(zhí)行這與原文檔的表述完全一致。Java SDK使用 Java 時(shí)必須在pom.xml中聲明對(duì) Direct Runner 的運(yùn)行時(shí)依賴dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-direct-java/artifactId version2.41.0/version scoperuntime/scope /dependency啟動(dòng)程序時(shí)通過args設(shè)置 Runner--runnerDirectRunner本倉庫的官方示例 WordCount.java 也明確說明通過--runnerYOUR_SELECTED_RUNNER即可切換 Runner使用 DirectRunner 時(shí)輸出為本地文件。Python SDK在 Python SDK 中默認(rèn) Runner 同樣是DirectRunner。運(yùn)行 wordcount 示例python -m apache_beam.examples.wordcount --input YOUR_INPUT_FILE --output counts對(duì)應(yīng)示例源碼位于 sdks/python/apache_beam/examples/wordcount.py。其底層實(shí)現(xiàn)在 sdks/python/apache_beam/runners/direct/direct_runner.py 中其中SwitchingDirectRunner會(huì)在FnApiRunner批處理吞吐高與BundleBasedDirectRunner支持流式執(zhí)行與某些尚未實(shí)現(xiàn)的原語之間自動(dòng)切換這從側(cè)面印證了 Python Direct Runner 在批/流場景下的覆蓋策略。3. Google Cloud Dataflow Runner托管式的云上執(zhí)行3.1 工作原理Google Cloud Dataflow Runner 使用 Cloud Dataflow 托管服務(wù)。運(yùn)行管線時(shí)Runner 會(huì)把你的可執(zhí)行代碼與依賴上傳到一個(gè) Google Cloud StorageGCS桶并創(chuàng)建 Cloud Dataflow 任務(wù)Job由該任務(wù)在 Google Cloud Platform 的托管資源上執(zhí)行管線。從 DataflowRunner.java 的源碼可以確認(rèn)其內(nèi)部大量操作圍繞 GCP 資源展開如上傳代碼、依賴檢測、與 Dataflow API 交互等。Dataflow Runner 與服務(wù)適用于大規(guī)模、持續(xù)運(yùn)行的作業(yè)并提供完全托管的服務(wù)fully managed service在整個(gè)作業(yè)生命周期內(nèi)自動(dòng)伸縮 worker 數(shù)量autoscaling動(dòng)態(tài)工作再平衡dynamic work rebalancing3.2 各 SDK 運(yùn)行方式Go SDK$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner dataflow \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latestJava SDK在pom.xml中聲明依賴dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-google-cloud-dataflow-java/artifactId version2.42.0/version scoperuntime/scope /dependency然后在 Maven JAR 插件中配置mainClassplugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-jar-plugin/artifactId version${maven-jar-plugin.version}/version configuration archive manifest addClasspathtrue/addClasspath classpathPrefixlib//classpathPrefix mainClassYOUR_MAIN_CLASS_NAME/mainClass /manifest /archive /configuration /plugin命令行運(yùn)行java -jar target/beam-examples-bundled-1.0.0.jar \ --runnerDataflowRunner \ --projectYOUR_GCP_PROJECT_ID \ --regionGCP_REGION \ --tempLocationgs://YOUR_GCS_BUCKET/temp/Python SDK先安裝 GCP 相關(guān)的擴(kuò)展組件pip install apache-beam[gcp]再運(yùn)行python -m apache_beam.examples.wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://YOUR_GCS_BUCKET/counts \ --runner DataflowRunner \ --project YOUR_GCP_PROJECT \ --region YOUR_GCP_REGION \ --temp_location gs://YOUR_GCS_BUCKET/tmp/4. Apache Flink Runner流式優(yōu)先的高吞吐執(zhí)行4.1 特性與適用場景Apache Flink Runner 基于 Apache Flink 執(zhí)行 Beam 管線既可以采用集群執(zhí)行模式如 YARN / Kubernetes / Mesos也可以采用本地嵌入式模式便于測試管線。Flink Runner 與 Flink 適合大規(guī)模、持續(xù)運(yùn)行的作業(yè)并提供流式優(yōu)先streaming-first的運(yùn)行時(shí)同時(shí)支持批處理與流數(shù)據(jù)處理程序同時(shí)支持極高吞吐與低事件延遲的運(yùn)行時(shí)具備 exactly-once 處理保證的容錯(cuò)能力流式程序中天然的反壓back-pressure機(jī)制自定義內(nèi)存管理可在內(nèi)存內(nèi)與內(nèi)存外數(shù)據(jù)處理算法之間高效、穩(wěn)健地切換與 YARN 以及 Apache Hadoop 生態(tài)其他組件的集成4.2 各 SDK 運(yùn)行方式Go SDK需要引入 flink runner 包并通過--endpoint指定 Runner 所在端點(diǎn)github.com/apache/beam/sdks/v2/go/pkg/beam/runners/flink$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner flink \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latest \ --endpointlocalhost:8081Java SDKPortable 模式從 Beam 2.18.0 起Docker Hub 提供了預(yù)構(gòu)建的 Flink Job Service 鏡像Flink 1.10、Flink 1.11、Flink 1.12、Flink 1.13、Flink 1.14啟動(dòng) JobService 端點(diǎn)docker run --nethost apache/beam_flink1.10_job_server:latest使用 PortableRunner 將管線提交到上述端點(diǎn)job_endpoint設(shè)為localhost:8099JobService 的默認(rèn)地址可選設(shè)置environment_type為LOOPBACK。示例import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions options PipelineOptions([ --runnerPortableRunner, --job_endpointlocalhost:8099, --environment_typeLOOPBACK ]) with beam.Pipeline(options) as p: ...Java SDK非 Portable 模式在pom.xml中聲明依賴dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-flink-1.14/artifactId version2.42.0/version /dependency命令行運(yùn)行mvn exec:java -Dexec.mainClassorg.apache.beam.examples.WordCount \ -Pflink-runner \ -Dexec.args--runnerFlinkRunner \ --inputFile/path/to/pom.xml \ --output/path/to/counts \ --flinkMasterflink master url \ --filesToStagetarget/word-count-beam-bundled-0.1.jarPython SDK與上述 Portable 模式完全一致先docker run --nethost apache/beam_flink1.10_job_server:latest啟動(dòng) JobService再通過PortableRunner、job_endpointlocalhost:8099、environment_typeLOOPBACK提交管線。5. Apache Spark Runner與 Spark 生態(tài)深度集成5.1 特性Apache Spark Runner 基于 Apache Spark 執(zhí)行 Beam 管線可以像原生 Spark 應(yīng)用一樣運(yùn)行以自包含應(yīng)用部署到本地模式、Spark Standalone 資源管理器或運(yùn)行在 YARN / Mesos 之上。Spark Runner 提供批處理、流式處理以及批流一體的管線與 RDD 和 DStream 相同的容錯(cuò)保證Spark 提供的安全特性基于 Spark 指標(biāo)系統(tǒng)的內(nèi)建指標(biāo)上報(bào)同時(shí)上報(bào) Beam Aggregators通過 Spark 的 Broadcast 變量原生支持 Beam side-inputs5.2 各 SDK 運(yùn)行方式Go SDK需要引入 spark runner 包并通過--endpoint指定端點(diǎn)github.com/apache/beam/sdks/v2/go/pkg/beam/runners/spark$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner spark \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latest \ --endpointlocalhost:8081Java SDK非 Portable 模式啟動(dòng) JobService 端點(diǎn)二選一使用 Docker推薦docker run --nethost apache/beam_spark_job_server:latest或從 Beam 源碼構(gòu)建./gradlew :runners:spark:3:job-server:runShadow通過 PortableRunner 提交 Python 管線到上述端點(diǎn)job_endpoint設(shè)為localhost:8099environment_type設(shè)為LOOPBACKPython 示例代碼同上文 Flink 部分。在pom.xml中聲明 Spark Runner 依賴dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-spark-3/artifactId version2.42.0/version /dependency并使用 Maven Shade 插件對(duì)應(yīng)用 JAR 做 shading避免打包沖突plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId configuration createDependencyReducedPomfalse/createDependencyReducedPom filters filter artifact*:*/artifact excludes excludeMETA-INF/*.SF/exclude excludeMETA-INF/*.DSA/exclude excludeMETA-INF/*.RSA/exclude /excludes /filter /filters /configuration executions execution phasepackage/phase goals goalshade/goal /goals configuration shadedArtifactAttachedtrue/shadedArtifactAttached shadedClassifierNameshaded/shadedClassifierName transformers transformer implementationorg.apache.maven.plugins.shade.resource.ServicesResourceTransformer/ /transformers /configuration /execution /executions /plugin命令行運(yùn)行mvn compile exec:java -Dexec.mainClassorg.apache.beam.examples.WordCount \ -Dexec.args--runnerSparkRunner --inputFilepom.xml --outputcounts -Pspark-runnerPython SDKpython -m apache_beam.examples.wordcount --input /path/to/inputfile \ --output /path/to/write/counts \ --runner SparkRunner從源碼層面看SparkRunner.java 提供了create()、create(SparkPipelineOptions)與fromOptions(PipelineOptions)等多種構(gòu)造入口并返回SparkPipelineResult表明 Beam 管線在 Spark 上以原生 Spark 應(yīng)用的形式被執(zhí)行與追蹤。6. Apache Samza Runner大規(guī)模有狀態(tài)流式作業(yè)6.1 特性Apache Samza Runner 在 Samza 應(yīng)用中執(zhí)行 Beam 管線可本地運(yùn)行也可以把應(yīng)用打包成.tgz部署到 YARN 集群或帶 Zookeeper 的 Samza standalone 集群。Samza Runner 與 Samza 適合大規(guī)模、有狀態(tài)的流式作業(yè)并提供對(duì)本地狀態(tài)基于 RocksDB 存儲(chǔ)的一流支持便于高頻流式作業(yè)快速訪問狀態(tài)支持狀態(tài)增量 checkpoint 而非全量快照的容錯(cuò)能力使 Samza 能擴(kuò)展到狀態(tài)極大的應(yīng)用完全異步的處理引擎讓遠(yuǎn)程調(diào)用更高效靈活的部署模型可在任何帶 Zookeeper 的托管環(huán)境中運(yùn)行應(yīng)用金絲雀發(fā)布canaries、升級(jí)與回滾等特性支持以最小停機(jī)時(shí)間支撐超大規(guī)模部署原文檔中的 Samza 小節(jié)僅面向 Java 與 Go SDK 展開。6.2 各 SDK 運(yùn)行方式Go SDK需要引入 samza runner 包并通過--endpoint指定端點(diǎn)github.com/apache/beam/sdks/v2/go/pkg/beam/runners/samza$ go install github.com/apache/beam/sdks/v2/go/examples/wordcount # As part of the initial setup, for non linux users - install package unix before run $ go get -u golang.org/x/sys/unix $ wordcount --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://your-gcs-bucket/counts \ --runner samza \ --project your-gcp-project \ --region your-gcp-region \ --temp_location gs://your-gcs-bucket/tmp/ \ --staging_location gs://your-gcs-bucket/binaries/ \ --worker_harness_container_imageapache/beam_go_sdk:latest \ --endpointlocalhost:8081Java SDK在pom.xml中聲明 Samza Runner 及其依賴dependency groupIdorg.apache.beam/groupId artifactIdbeam-runners-samza/artifactId version2.42.0/version scoperuntime/scope /dependency !-- Samza dependencies -- dependency groupIdorg.apache.samza/groupId artifactIdsamza-api/artifactId version${samza.version}/version /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-core_2.11/artifactId version${samza.version}/version /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-kafka_2.11/artifactId version${samza.version}/version scoperuntime/scope /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-kv_2.11/artifactId version${samza.version}/version scoperuntime/scope /dependency dependency groupIdorg.apache.samza/groupId artifactIdsamza-kv-rocksdb_2.11/artifactId version${samza.version}/version scoperuntime/scope /dependency命令行運(yùn)行$ mvn exec:java -Dexec.mainClassorg.apache.beam.examples.WordCount \ -Psamza-runner \ -Dexec.args--runnerSamzaRunner \ --inputFile/path/to/input \ --output/path/to/counts7. Apache Nemo Runner編譯期優(yōu)化與分布式執(zhí)行7.1 特性Apache Nemo Runner 基于 Apache Nemo 執(zhí)行 Beam 管線通過 Nemo 編譯器中的多種優(yōu)化 passoptimization passes對(duì) Beam 管線做優(yōu)化再在 Nemo 運(yùn)行時(shí)上分布式執(zhí)行。你也可以把應(yīng)用作為自包含程序部署到本地模式或使用 YARN、Mesos 等資源管理器運(yùn)行。Nemo Runner 提供批處理與流式管線容錯(cuò)能力與 YARN 以及 Apache Hadoop 生態(tài)其他組件的集成對(duì) Nemo 優(yōu)化器提供的多種優(yōu)化的支持7.2 Java 運(yùn)行方式在pom.xml中聲明依賴注意需排除 slf4j 沖突dependency groupIdorg.apache.nemo/groupId artifactIdnemo-compiler-frontend-beam/artifactId version${nemo.version}/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-common/artifactId version${hadoop.version}/version exclusions exclusion groupIdorg.slf4j/groupId artifactIdslf4j-api/artifactId /exclusion exclusion groupIdorg.slf4j/groupId artifactIdslf4j-log4j12/artifactId /exclusion /exclusions /dependency原文檔建議自包含應(yīng)用可能更易于管理并能完整使用 Nemo 提供的功能。只需添加上述依賴并用 Maven Shade 插件對(duì)應(yīng)用 JAR 做 shadingShade 配置與第 5.2 節(jié) Spark 的完全相同可復(fù)用。命令行運(yùn)行$ mvn package -Pnemo-runner java -cp target/word-count-beam-bundled-0.1.jar org.apache.beam.examples.WordCount \ --runnerNemoRunner --inputFilepwd/pom.xml --outputcounts8. Hazelcast Jet Runner實(shí)驗(yàn)性的高吞吐內(nèi)存計(jì)算8.1 特性與現(xiàn)狀說明Hazelcast Jet Runner 基于 Hazelcast Jet 執(zhí)行 Beam 管線適合大規(guī)模持續(xù)作業(yè)并提供同時(shí)支持批處理有界與流式無界數(shù)據(jù)集同時(shí)支持極高吞吐與低事件延遲的運(yùn)行時(shí)流式程序中天然的反壓機(jī)制帶內(nèi)存存儲(chǔ)的分布式大規(guī)模并行處理引擎需要特別注意Jet Runner 目前處于 EXPERIMENTAL實(shí)驗(yàn)性狀態(tài)尚不能使用 Jet 的許多能力Jet 本身具備完整的容錯(cuò)支持但 Jet Runner沒有作業(yè)失敗后必須重啟Jet 內(nèi)部性能極高但 Runner 目前無法匹敵因?yàn)?Beam 管線的優(yōu)化/改造surgery尚未完全實(shí)現(xiàn)。8.2 Java 運(yùn)行方式$ mvn package -P jet-runner java -cp target/word-count-beam-bundled-0.1.jar org.apache.beam.examples.WordCount \ --runnerJetRunner --jetLocalMode3 --inputFilepwd/pom.xml --outputcounts其中--jetLocalMode3用于指定 Jet 的本地運(yùn)行模式。9. Runner 選型速查與實(shí)戰(zhàn)建議結(jié)合原文檔與倉庫源碼可以把上述 Runner 歸納為三類使用場景Runner場景定位關(guān)鍵運(yùn)行參數(shù)備注Direct Runner本地開發(fā)、單元測試、模型語義校驗(yàn)--runnerDirectRunnerJava/ 默認(rèn)Go、Python默認(rèn)開啟不可變性與可編碼性檢查Dataflow Runner云上大規(guī)模持續(xù)作業(yè)--runnerDataflowRunner--project/--region/--tempLocation托管服務(wù)、自動(dòng)伸縮、動(dòng)態(tài)再平衡Flink Runner流式優(yōu)先的高吞吐批流作業(yè)--runnerFlinkRunner/ PortableRunner --job_endpoint支持 YARN / Kubernetes / 本地嵌入式Spark Runner與 Spark 生態(tài)深度集成的批流作業(yè)--runnerSparkRunner-Pspark-runner側(cè)輸入走 Broadcast 變量Samza Runner大規(guī)模有狀態(tài)流式作業(yè)--runnerSamzaRunner-Psamza-runnerRocksDB 本地狀態(tài)、增量 checkpointNemo Runner追求編譯期優(yōu)化的分布式執(zhí)行--runnerNemoRunner-Pnemo-runner需 shade 應(yīng)用 JARJet Runner實(shí)驗(yàn)性高吞吐內(nèi)存計(jì)算--runnerJetRunner --jetLocalMode3實(shí)驗(yàn)狀態(tài)作業(yè)失敗需重啟選型建議本地開發(fā)與測試首選 Direct Runner。它能提前發(fā)現(xiàn)元素被篡改、Coder 不匹配、依賴非模型保證語義等問題且可以直接使用本地調(diào)試工具需要托管服務(wù)與彈性伸縮時(shí)選擇 Dataflow適合云上大規(guī)模持續(xù)作業(yè)流式優(yōu)先、需要 exactly-once 與反壓支持時(shí)選擇 Flink與 Spark 生態(tài)RDD/DStream、Broadcast、Spark 指標(biāo)深度集成時(shí)選擇 Spark大規(guī)模有狀態(tài)流式作業(yè)選擇 SamzaRocksDB 本地狀態(tài) 增量 checkpoint實(shí)驗(yàn)性探索可選擇 Nemo 或 Jet但要接受 Jet Runner 當(dāng)前不支持完整容錯(cuò)、作業(yè)失敗需重啟的現(xiàn)實(shí)。需要補(bǔ)充說明的是本文涉及的依賴版本號(hào)如beam-runners-direct-java2.41.0、beam-runners-google-cloud-dataflow-java2.42.0 等均來自原文檔實(shí)際使用時(shí)應(yīng)以你當(dāng)前構(gòu)建工具解析到的最新版本為準(zhǔn)${samza.version}、${nemo.version}、${hadoop.version}等占位符需在你的構(gòu)建中顯式定義。此外Dataflow 運(yùn)行還需要預(yù)先配置 GCP 項(xiàng)目、GCS 桶等環(huán)境資源本倉庫的 runners/google-cloud-dataflow-java 與 runners/spark、runners/flink 等目錄提供了各 Runner 的完整實(shí)現(xiàn)可作為進(jìn)一步研究底層調(diào)度與翻譯邏輯的入口。贊分享【免費(fèi)下載鏈接】beamApache Beam is a unified programming model for Batch and Streaming data processing.項(xiàng)目地址https://gitcode.com/gh_mirrors/beam18/beam點(diǎn)擊查看免費(fèi)下載相關(guān)推薦Apache Beam Runner 概念詳解從 DirectRunner 到分布式 Runner 的選擇與實(shí)戰(zhàn)Apache Beam Runner 概念詳解從 DirectRunner 到分布式 Runner 的選擇與實(shí)戰(zhàn) Apache Beam 的核心價(jià)值在于一次大數(shù)據(jù)批處理流處理數(shù)據(jù)工程Warp 補(bǔ)全引擎的 Basic Parser 架構(gòu)從 Lex 到類型驅(qū)動(dòng) Full Parse 的遞歸下降解析器全解析Warp 補(bǔ)全引擎的 Basic Parser 架構(gòu)從 Lex 到類型驅(qū)動(dòng) Full Parse 的遞歸下降解析器全解析 導(dǎo)讀 本文深入剖析 Warp 開源倉批處理流處理大數(shù)據(jù)Apache Beam Direct Runner 完全指南在本地運(yùn)行與調(diào)試 Beam 管道Apache Beam Direct Runner 完全指南在本地運(yùn)行與調(diào)試 Beam 管道 導(dǎo)讀 Apache Beam 提供了統(tǒng)一的批處理與流處理編程模型大數(shù)據(jù)批處理流處理數(shù)據(jù)工程創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考