Hi,
I'm working on a custom TimestampAssigner which will do different things
depending on the value of the extracted timestamp. One of the actions I
want to take is to drop messages entirely if their timestamp meets certain
criteria.
Of course there's no direct way to do this in the
Hi all,
I'm looking into what happens when messages are ingested with timestamps
far into the future (e.g. due to corruption or a wrong clock at the sender).
I'm aware of the effect on watermarking, but another thing I'm concerned
about is the performance impact of the extra windows this will
)
> java_execution_env.addDefaultKryoSerializer(Message,
> MessageKryoSerializer)
>
> With or without these lines, the job crashes with a KryoException (full
> stack trace at https://pastebin.com/zxxzCqH0), it doesn't appear that
> addDefaultKryoSerializer is doing anything.
>
Hi,
Is there any way to execute a job using the LocalEnvironment when using the
Python streaming API? This would make it much easier to debug jobs.
At the moment I'm not aware of any way of running them except firing up a
local cluster and submitting the job with pyflink-stream.sh.
Thanks,
Joe
with a KryoException (full
stack trace at https://pastebin.com/zxxzCqH0), it doesn't appear that
addDefaultKryoSerializer is doing anything.
Is there an officially supported way to set custom serializers in Python?
Thanks,
Joe Malt
Engineering Intern, Stream Processing
Yelp
for
fc:org.apache.flink.runtime.operators.testutils.MockEnvironmentBuilder
Best,
Joe Malt
Engineering Intern, Stream Processing
Yelp Inc.
On Fri, Aug 24, 2018 at 4:50 AM, Chang Liu wrote:
> No worries, I found it here:
>
>
> org.apache.flink
> flink-runtime_${scala.binary.version}
> ${flink.versi
in Jython.
I've tried adding a MapFunction that maps each input to String(input)where
String is the constructor for java.lang.String. This made no difference; I
get the same error.
Any ideas?
Thanks,
Joe Malt
Software Engineering Intern
Yelp
fine to me.
>
> An example of an official document may not guarantee your success due to
> maintenance issues.
>
> cc @Chesnay
>
> [1]: https://github.com/apache/flink/blob/master/flink-libraries/flink-
> streaming-python/src/test/python/org/apache/flink/
> streaming/python
ciated.
Thanks,
Joe Malt
Software Engineer Intern
Yelp
s of Flink can't find the class for some
reason.
The command I'm using to run the pipeline:
./pyflink-stream.sh /Users/jmalt/flink-python/KafkaRead.py
/Users/jmalt/flink-python/MyCustomKafkaDeserializer.py - --local
How can I make Flink see the custom deserializer?
Thanks,
Joe Malt
Software E
nStreamBinder -v
"$FLINK_ROOT_DIR"/opt/flink-streaming-python*.jar "$@"
Thanks,
Joe Malt
Engineering Intern, Stream Processing
Yelp
11 matches
Mail list logo