我有一个用 Go 编写的 Beam 批处理管道,它需要一个 2000 万行的 .csv 文件(大约 600 MB 的数据),执行基本的转换步骤,例如 SumPerKey 并将输出写回 GCS。
在 Dataflow 上运行管道时,它仅调用一个包含 1 个运行器的池!
我原以为 Dataflow 会针对这种数据量在多个工作人员之间并行处理作业。我错过了什么吗?
这是我的代码:
func main() {
flag.Parse()
beam.Init()
p, s := beam.NewPipelineWithRoot()
ctx := context.Background()
log.Infof(ctx, "Started pipeline on scope: %s", s)
/* [TEST PIPELINE START ]*/
sr := csvio.Read(s, *input, reflect.TypeOf(Rating{}))
pwo := beam.ParDo(s.Scope("Pair Key With One"),
func(x Rating, emit func(int, int)) {
emit(x.UserId, 1)
}, sr)
spk := stats.SumPerKey(s, pwo)
mp := beam.ParDo(s.Scope("Map KV To Struct"),
func(k int, v int, emit func(UserRatings)) {
emit(UserRatings{
UserId: k,
Ratings: v,
})
}, spk)
t := top.Largest(s, mp, 1000, func(x, y UserRatings) bool { return x.Ratings < y.Ratings })
o := beam.ParDo(s, func(x []UserRatings) string {
if data, err := json.MarshalIndent(x, "", ""); err != nil {
return fmt.Sprintf("[Err]: %v", err)
} else {
return fmt.Sprintf("Output: %s", data)
}
}, t)
textio.Write(s, *output, o)
/* [TEST PIPELINE END ]*/
if err := beamx.Run(ctx, p); err != nil {
fmt.Println(err)
log.Exitf(ctx, "Failed to execute job: on ctx=%v:")
}
}
我通过此命令行部署管道:
go run main.go \
--runner dataflow \
--max_num_workers 10 \
--file gs://${BUCKET?}/ratings.csv \
--output gs://${BUCKET?}/reporting.txt \
--project ${PROJECT?} \
--temp_location gs://${BUCKET?}/tmp/ \
--staging_location gs://${BUCKET?}/binaries/ \
--worker_harness_container_image=gcr.io/drawndom-app/beam/go:latest
注意:当我将 --num_workers 设置为 5 时,它会调用 5 个 worker ,但我希望它自动执行此操作。
最佳答案
更新:
由于这个 lib,我在 .csv 输入之前添加了 Reshuffle 步骤Dataflow 能够通过再添加 1 个工作器来进行自动缩放。
我仍然需要了解如何优化管道的并行性。
使用的代码:
func Reshuffle(s beam.Scope, col beam.PCollection) beam.PCollection {
s = s.Scope("Reshuffle")
col = beam.ParDo(s, func(x beam.X) (int, beam.X) {
return rand.Int(), x
}, col)
col = beam.GroupByKey(s, col)
return beam.ParDo(s, func(key int, values func(*beam.X) bool, emit func(beam.X)) {
var x beam.X
for values(&x) {
emit(x)
}
}, col)
}
关于go - Apache Beam Go SDK - 数据流无法正确自动缩放(并行化步骤),我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/57922410/
我在从html页面生成PDF时遇到问题。我正在使用PDFkit。在安装它的过程中,我注意到我需要wkhtmltopdf。所以我也安装了它。我做了PDFkit的文档所说的一切......现在我在尝试加载PDF时遇到了这个错误。这里是错误:commandfailed:"/usr/local/bin/wkhtmltopdf""--margin-right""0.75in""--page-size""Letter""--margin-top""0.75in""--margin-bottom""0.75in""--encoding""UTF-8""--margin-left""0.75in""-
我主要使用Ruby来执行此操作,但到目前为止我的攻击计划如下:使用gemsrdf、rdf-rdfa和rdf-microdata或mida来解析给定任何URI的数据。我认为最好映射到像schema.org这样的统一模式,例如使用这个yaml文件,它试图描述数据词汇表和opengraph到schema.org之间的转换:#SchemaXtoschema.orgconversion#data-vocabularyDV:name:namestreet-address:streetAddressregion:addressRegionlocality:addressLocalityphoto:i
我对最新版本的Rails有疑问。我创建了一个新应用程序(railsnewMyProject),但我没有脚本/生成,只有脚本/rails,当我输入ruby./script/railsgeneratepluginmy_plugin"Couldnotfindgeneratorplugin.".你知道如何生成插件模板吗?没有这个命令可以创建插件吗?PS:我正在使用Rails3.2.1和ruby1.8.7[universal-darwin11.0] 最佳答案 随着Rails3.2.0的发布,插件生成器已经被移除。查看变更日志here.现在
我尝试运行2.x应用程序。我使用rvm并为此应用程序设置其他版本的ruby:$rvmuseree-1.8.7-head我尝试运行服务器,然后出现很多错误:$script/serverNOTE:Gem.source_indexisdeprecated,useSpecification.Itwillberemovedonorafter2011-11-01.Gem.source_indexcalledfrom/Users/serg/rails_projects_terminal/work_proj/spohelp/config/../vendor/rails/railties/lib/r
我正在查看instance_variable_set的文档并看到给出的示例代码是这样做的:obj.instance_variable_set(:@instnc_var,"valuefortheinstancevariable")然后允许您在类的任何实例方法中以@instnc_var的形式访问该变量。我想知道为什么在@instnc_var之前需要一个冒号:。冒号有什么作用? 最佳答案 我的第一直觉是告诉你不要使用instance_variable_set除非你真的知道你用它做什么。它本质上是一种元编程工具或绕过实例变量可见性的黑客攻击
我正在尝试在我的centos服务器上安装therubyracer,但遇到了麻烦。$geminstalltherubyracerBuildingnativeextensions.Thiscouldtakeawhile...ERROR:Errorinstallingtherubyracer:ERROR:Failedtobuildgemnativeextension./usr/local/rvm/rubies/ruby-1.9.3-p125/bin/rubyextconf.rbcheckingformain()in-lpthread...yescheckingforv8.h...no***e
我花了三天的时间用头撞墙,试图弄清楚为什么简单的“rake”不能通过我的规范文件。如果您遇到这种情况:任何文件夹路径中都不要有空格!。严重地。事实上,从现在开始,您命名的任何内容都没有空格。这是我的控制台输出:(在/Users/*****/Desktop/LearningRuby/learn_ruby)$rake/Users/*******/Desktop/LearningRuby/learn_ruby/00_hello/hello_spec.rb:116:in`require':cannotloadsuchfile--hello(LoadError) 最佳
我在pry中定义了一个函数:to_s,但我无法调用它。这个方法去哪里了,怎么调用?pry(main)>defto_spry(main)*'hello'pry(main)*endpry(main)>to_s=>"main"我的ruby版本是2.1.2看了一些答案和搜索后,我认为我得到了正确的答案:这个方法用在什么地方?在irb或pry中定义方法时,会转到Object.instance_methods[1]pry(main)>defto_s[1]pry(main)*'hello'[1]pry(main)*end=>:to_s[2]pry(main)>defhello[2]pry(main)
我使用的是Firefox版本36.0.1和Selenium-Webdrivergem版本2.45.0。我能够创建Firefox实例,但无法使用脚本继续进行进一步的操作无法在60秒内获得稳定的Firefox连接(127.0.0.1:7055)错误。有人能帮帮我吗? 最佳答案 我遇到了同样的问题。降级到firefoxv33后一切正常。您可以找到旧版本here 关于ruby-无法在60秒内获得稳定的Firefox连接(127.0.0.1:7055),我们在StackOverflow上找到一个类
当我尝试安装Ruby时遇到此错误。我试过查看this和this但无济于事➜~brewinstallrubyWarning:YouareusingOSX10.12.Wedonotprovidesupportforthispre-releaseversion.Youmayencounterbuildfailuresorotherbreakages.Pleasecreatepull-requestsinsteadoffilingissues.==>Installingdependenciesforruby:readline,libyaml,makedepend==>Installingrub