java – Apache Beam中有状态处理的问题
内容导读
互联网集市收集整理的这篇技术教程文章主要介绍了java – Apache Beam中有状态处理的问题,小编现在分享给大家,供广大互联网技能从业者学习和参考。文章包含1888字,纯文字阅读大概需要3分钟。
内容图文
所以我读过梁的stateful processing和timely processing文章,并发现了实现这些功能的问题.
我试图解决的问题类似于this,为每一行生成一个顺序索引.因为我希望能够将数据流生成的行引用到原始源的行.
public static class createIndex extends DoFn<String, KV<String, String>> {
@StateId("count")
private final StateSpec<ValueState<Long>> countState = StateSpecs.value(VarLongCoder.of());
@ProcessElement
public void processElement(ProcessContext c, @StateId("count") ValueState<Long> countState) {
String val = c.element();
long count = 0L;
if(countState.read() != null)
count = countState.read();
count = count + 1;
countState.write(count);
c.output(KV.of(String.valueOf(count), val));
}
}
Pipeline p = Pipeline.create(options);
p.apply(TextIO.read().from("gs://randomBucket/file.txt"))
.apply(ParDo.of(new createIndex()));
我按照我在网上找到的任何内容查看了ParDo的原始源代码,但不确定需要做什么.我得到的错误是:
java.lang.IllegalArgumentException: ParDo requires its input to use KvCoder in order to use state and timers.
我意识到这是一个简单的问题,但由于缺乏足够的示例或文档,我无法解决问题.我很感激任何帮助.谢谢!
解决方法:
好的,所以我继续研究这个问题并阅读一些源代码,并且能够解决问题.事实证明,ParDo.of(new DoFn())的输入要求输入的输入为KV< T,T>.因此,为了读取文件并为每一行创建索引,我需要通过Key Value Pair对象传递它.下面我添加了代码:
public static class FakeKvPair extends DoFn<String, KV<String, String>> {
@ProcessElement
public void processElement(ProcessContext c) {
c.output(KV.of("", c.element()));
}
}
并将管道更改为:
Pipeline p = Pipeline.create(options);
p.apply(TextIO.read().from("gs://randomBucket/file.txt"))
.apply(ParDo.of(new FakeKvPair()))
.apply(ParDo.of(new createIndex()));
出现的新问题是,在我的计数中是否保留了行的顺序,因为我正在运行一个额外的ParDo函数,该函数可能会改变输入到createIndex()的行的顺序.
在我的本地机器上保留订单,但我不知道它将如何扩展到Dataflow.但我会问这是一个不同的问题.
内容总结
以上是互联网集市为您收集整理的java – Apache Beam中有状态处理的问题全部内容,希望文章能够帮你解决java – Apache Beam中有状态处理的问题所遇到的程序开发问题。 如果觉得互联网集市技术教程内容还不错,欢迎将互联网集市网站推荐给程序员好友。
内容备注
版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 gblab@vip.qq.com 举报,一经查实,本站将立刻删除。
内容手机端
扫描二维码推送至手机访问。