You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
This repository has been archived by the owner on Dec 14, 2022. It is now read-only.
Describe the bug
When using a Flink DeserializationSchema that has overridden the isEndOfStream method, the source does not terminate when isEndOfStream evaluates to true.
To Reproduce
Steps to reproduce the behavior:
Create a DeserializationSchema that returns true for isEndOfStream
Use FlinkPulsarSource with the schema
Stream continues indefinitely
Expected behavior
The FlinkPulsarSource closes with a maximum timestamp to flush downstream nodes.
Additional context
I am currently attempting to use this in a very simple flink job to export event-time-bounded selections data from pulsar into an S3 bucket. The job will not finalize the file till the stream terminates. If I manually cancel the stream the data shows up in s3.
The text was updated successfully, but these errors were encountered:
Describe the bug
When using a Flink DeserializationSchema that has overridden the
isEndOfStream
method, the source does not terminate whenisEndOfStream
evaluates to true.To Reproduce
Steps to reproduce the behavior:
isEndOfStream
Expected behavior
The FlinkPulsarSource closes with a maximum timestamp to flush downstream nodes.
Additional context
I am currently attempting to use this in a very simple flink job to export event-time-bounded selections data from pulsar into an S3 bucket. The job will not finalize the file till the stream terminates. If I manually cancel the stream the data shows up in s3.
The text was updated successfully, but these errors were encountered: