【问题标题】:Scheduler closing stream warning调度程序关闭流警告
【发布时间】:2018-11-27 21:13:08
【问题描述】:

我的笔记本电脑上有一个定期批处理作业。代码如下所示:

client = Client()
print(client.scheduler_info())
topic='raw_data'
start = datetime.datetime.now()
delta = datetime.timedelta(minutes=2)
while True:
    end = start + delta
    if end <= datetime.datetime.now():
        start = end
        print('It\'s time to run the analysis for the 2 mins')
        data = get_data_from_parquet('raw_data_fast_par.par', start=start, end=end)
        metrics = [Metric1(), Metric2(), Metric3()]
        print(data.npartitions)
        channels = data.groupby(['col1', 'col2', 'col3'])
        for metric in metrics:
            features = metric.map_job(channels, start, end)
            print(features.count().compute())

简而言之,我每两分钟对数据进行某种分析,我从 parquet 文件中读取这些数据,预测日期过滤。这是一个测试,所以我知道现在它没有多大意义。 我在终端上收到以下警告。有人可以解释为什么会发生这种情况,如果它很重要,以及如何避免它?

distributed.comm.tcp - WARNING - Closing dangling stream in <TCP local=tcp://127.0.0.1:55448 remote=tcp://127.0.0.1:42197>

【问题讨论】:

    标签: python dask dask-distributed


    【解决方案1】:

    我不知道实际问题是什么,但您可以尝试在完成后彻底关闭本地集群,或许可以使用 Client 作为上下文管理器。

    with Client() as client:
        ...
    

    【讨论】:

    • 将尝试返回。我确实尝试在 for 循环结束时重新启动,但没有任何影响,我仍然收到该消息
    • 这似乎已经摆脱了警告,但现在我失去了执行的“历史”数据。有没有这方面的工作?
    • 我设法解决了这个问题,首先实例化了一个本地集群,然后用cluster = LocalCluster() client=Client(cluster) 将其分配给客户端,尽管我不知道如何解释。
    • @Apostolos 记录了LocalCluster 的使用here。我使用with LocalCluster(**kwargs) as cluster: 和嵌套的with Client(cluster) as client: 来减轻悬空流 警告。我之前一直在使用cluster.close()client.close()with 语句似乎更可靠地工作。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-26
    • 2015-05-19
    • 2011-04-05
    • 2023-03-09
    相关资源
    最近更新 更多