【问题标题】:Pyflink windowAll() by event-time to apply a clutering modelFlink windowAll() 按事件时间应用聚类模型
【发布时间】:2022-07-29 23:30:52
【问题描述】:

我是 pyflink 框架的初学者,我想知道我的用例是否可以使用它...

我需要制作一个翻滚窗口并在其上应用 python udf(scikit 学习聚类模型)。 用例是:每 30 秒我想对前 30 秒的数据应用我的 udf。

目前我成功地在流中使用来自 kafka 的数据,但是我无法使用 python API 在非键控流上创建 30 秒窗口。

你知道我的用例的一些例子吗?你知道 pyflink API 是否允许这样做吗?

这是我的第一枪:

from pyflink.common import Row
from pyflink.common.serialization import JsonRowDeserializationSchema, JsonRowSerializationSchema
from pyflink.common.typeinfo import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.connectors import FlinkKafkaConsumer, FlinkKafkaProducer
from pyflink.common.watermark_strategy import TimestampAssigner, WatermarkStrategy
from pyflink.common import Duration

import time

from utils.selector import Selector
from utils.timestampAssigner import KafkaRowTimestampAssigner

# 1. create a StreamExecutionEnvironment
env = StreamExecutionEnvironment.get_execution_environment()
# the sql connector for kafka is used here as it's a fat jar and could avoid dependency issues
env.add_jars("file:///flink-sql-connector-kafka_2.11-1.14.0.jar")

deserialization_schema = JsonRowDeserializationSchema.builder() \
    .type_info(type_info=Types.ROW_NAMED(["labelId","freq","timestamp"],[Types.STRING(),Types.DOUBLE(),Types.STRING()])).build()


kafka_consumer = FlinkKafkaConsumer(
    topics='events',
    deserialization_schema=deserialization_schema,
    properties={'bootstrap.servers': 'localhost:9092'})



# watermark_strategy = WatermarkStrategy.for_bounded_out_of_orderness(Duration.of_seconds(5))\
#     .with_timestamp_assigner(KafkaRowTimestampAssigner())

ds = env.add_source(kafka_consumer)
ds.print()
ds = ds.windowAll()
# ds.print()

env.execute()


WARNING: An illegal reflective access operation has occurred
WARNING: Illegal reflective access by org.apache.flink.api.java.ClosureCleaner (file:/home/dorian/dataScience/pyflink/pyflink_env/lib/python3.6/site-packages/pyflink/lib/flink-dist_2.11-1.14.0.jar) to field java.util.Properties.serialVersionUID
WARNING: Please consider reporting this to the maintainers of org.apache.flink.api.java.ClosureCleaner
WARNING: Use --illegal-access=warn to enable warnings of further illegal reflective access operations
WARNING: All illegal access operations will be denied in a future release
Traceback (most recent call last):
  File "/home/dorian/dataScience/pyflink/project/__main__.py", line 35, in <module>
    ds = ds.windowAll()
AttributeError: 'DataStream' object has no attribute 'windowAll'

谢谢

【问题讨论】:

  • 您说您需要制作一个缩略图窗口,但您使用的是 .windowAll()。你这样做只是为了测试还是有混乱?您可以使用keyBy() 为您的流设置键控。
  • 嗯,是的,也许我误解了一些东西,我似乎可以在 non_stream 上制作一个缩略图窗口,所以使用 windowAll(),至少在 java 中:``` DataStream globalResults = resultsPerKey。 windowAll(TumblingEventTimeWindows.of(Time.seconds(5))) .process(new TopKWindowFunction()); ```
  • 是的,我认为您的示例 resultsPerKey .windowAll(TumblingEventTimeWindows.of(Time.seconds(5))) 应该可以工作。请注意,在非键控流的情况下,您的原始流不会被拆分为多个逻辑流,并且所有窗口逻辑都将由单个任务执行,即并行度为 1。如果您的流很大,您可以拥有一些性能问题。
  • 在 Pyflink 'AttributeError: 'DataStream' object has no attribute 'windowAll' 中似乎是不可能的。难道还没实现?
  • 我对 PyFlink 不熟悉。我在 python window assigners的文档中找不到任何提及 .windowAll()

标签: apache-flink flink-streaming pyflink flinkml


【解决方案1】:

我遇到了同样的错误/问题,你能解决这个问题吗?你也能帮帮我吗?

【讨论】:

猜你喜欢
  • 1970-01-01
  • 2018-01-17
  • 1970-01-01
  • 2018-11-05
  • 2019-06-10
  • 2018-11-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多