Using partition "breaks" program logic

Abierto
#427 2 comentarios 0 reacciones 0 asignados Ver en GitHub

Nadie ha tomado este issue todavía.

Evaluación

Dificultad
4/5
Tiempo estimado
3-5 días
Aptitud para principiantes
28/100
Tipo de issue
Error
Claridad
Bastante claro
Estado de actividad
Estancado
Stack tecnológico
python

Línea de trabajo

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.

Escrito por el modelo de indexación a partir del texto del issue.

Descripción

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).

Lenguaje dominante
Python
Estrellas
1.3k
Forks
149
Métricas de merge de PR
Sin PR fusionados en 30 d

Guía de contribución

Abrir la guía de contribución

Primeros pasos

  1. Lee el issue completo y luego la guía de contribución del proyecto.
  2. Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
  3. Haz un fork del repositorio y trabaja en una rama.
  4. Abre un pull request que haga referencia al número del issue.

Más de python-streamz/streamz

Todos los issues de python-streamz/streamz

Issues similares

Más issues de Python

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.