Using partition "breaks" program logic

オープン
#427 コメント 2 件 リアクション 0 件 担当者 0 名 GitHub で見る

まだ誰も着手していません。

評価

難易度
4/5
見積もり時間
3〜5日
初心者へのやさしさ
28/100
issue の種類
バグ
明瞭さ
おおむね明確
活発さ
停滞
技術スタック
python

調査の方向性

Start by tracing stream.partition, stream.emit, and accumulate, then read the documented async def process_file example in “Processing Time and Back Pressure.” Determine how a partition flushes when input reaches EOF and how callers can wait for pending processing; done means the final count includes all lines before the concluding print runs.

索引モデルが issue の本文から書いたものです。

説明

I am struggling to use partition in a pipeline because it "breaks" the logic of my program; presumably because it introduces asynchronous processing.

As a simplified example, I have something that works along the lines of this:

import streamz


def main():
    state = {
        "cnt": 0,
    }
    stream = streamz.Stream()
    cntd = stream.accumulate(cnt,
                             returns_state=True,
                             start=state)
    cntd.sink(print)

    with open("many_lines.txt", "r") as fh:
        for line in fh:
            stream.emit(line)
    print(f"found {state.get('cnt')} lines")


def cnt(state, itm):
    state["cnt"] += 1
    return state, itm


if __name__ == "__main__":
    main()

This basically runs through all the lines in the file many_lines.txt, counts and prints them and then reports

found 10000 lines

So far so good.

When I introduce partition now, like this:

import streamz


def main():
    state = {
        "cnt": 0,
    }
    stream = streamz.Stream()
    parted = stream.partition(10001, timeout=2)  # <= PARTITION HERE
    cntd = parted.accumulate(cnt,
                             returns_state=True,
                             start=state)
    cntd.sink(print)

    with open("many_lines.txt", "r") as fh:
        for line in fh:
            stream.emit(line)
    print(f"found {state.get('cnt')} lines")


def cnt(state, itm):
    state["cnt"] += 1
    return state, itm


if __name__ == "__main__":
    main()

I would want to see basically the same result. But I see nothing for some time and then

found 0 lines

I know, there are only 10'000 lines in many_lines.txt so the partition will never fill up, but it should hit the timeout at some point and "release" the data, no?

I suspect that the program terminates before the partition hits the timeout, so I tried (many variations of) awaiting stream.emit(line). That was inspired by the async def process_file(fn): function in Processing Time and Back Pressure.

For example like this:

import streamz


def main():
    state = {
        "cnt": 0,
    }
    stream = streamz.Stream()
    parted = stream.partition(10001, timeout=2)
    cntd = parted.accumulate(cnt,
                             returns_state=True,
                             start=state)
    cntd.sink(print)

    with open("many_lines.txt", "r") as fh:
        for line in fh:
            await stream.emit(line)  # <= USE AWAIT HERE
    print(f"found {state.get('cnt')} lines")


def cnt(state, itm):
    state["cnt"] += 1
    return state, itm


if __name__ == "__main__":
    main()

But this (obviously) does not work (SyntaxError: 'await' outside async function). And I also did not find a way to make it work.

(How) Can I make sure the for loop terminates before the print statement (or any remaining code, for that matter) is executed? Or am I getting this completely wrong?

My use case is to read (all) lines in pretty big files (I cannot load into memory at once), send them through a streamz pipeline and then continue with my program. "Then" meaning, after all lines are processed (also those that might be "stuck" in a partition when no more lines are emitted because we reached EOF; this is why I need the timeout, I believe).

主要言語
Python
スター
1.3k
フォーク
149
PR マージ指標
30日以内にマージされた PR はありません

コントリビューションガイド

コントリビューションガイドを開く

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

python-streamz/streamz のほかの issue

python-streamz/streamz の issue をすべて見る

似ている issue

Python の issue をもっと見る

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。