Hi,

I fail to run the PopularPlacesFromKafka example with the following
exception, and I wonder what might cause this "Invalid record" error?

when running within Intellij IDEA -->
07:52:23.960 [Source: Custom Source -> Map (7/8)] INFO
 org.apache.flink.runtime.taskmanager.Task  - Source: Custom Source -> Map
(7/8) (930e95aac65cbda39d9f1eaa41891253) switched from RUNNING to FAILED.
java.lang.RuntimeException: Invalid record:
4010,2013003778,2013003775,START,2013-01-01 00:13:00,1970-01-01
00:00:00,-74.00074,40.7359,-73.98559,40.739063,1
at
com.dataartisans.flinktraining.exercises.datastream_java.datatypes.TaxiRide.fromString(TaxiRide.java:119)
~[flink-training-exercises-0.15.1.jar:na]
at
com.dataartisans.flinktraining.exercises.datastream_java.utils.TaxiRideSchema.deserialize(TaxiRideSchema.java:37)
~[flink-training-exercises-0.15.1.jar:na]
at
com.dataartisans.flinktraining.exercises.datastream_java.utils.TaxiRideSchema.deserialize(TaxiRideSchema.java:28)
~[flink-training-exercises-0.15.1.jar:na]
at
org.apache.flink.streaming.util.serialization.KeyedDeserializationSchemaWrapper.deserialize(KeyedDeserializationSchemaWrapper.java:42)
~[flink-training-exercises-0.15.1.jar:na]
at
org.apache.flink.streaming.connectors.kafka.internal.Kafka09Fetcher.runFetchLoop(Kafka09Fetcher.java:139)
~[flink-training-exercises-0.15.1.jar:na]
at
org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase.run(FlinkKafkaConsumerBase.java:652)
~[flink-training-exercises-0.15.1.jar:na]
at
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:86)
~[flink-streaming-java_2.11-1.4.2.jar:1.4.2]
at
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:55)
~[flink-streaming-java_2.11-1.4.2.jar:1.4.2]
at
org.apache.flink.streaming.runtime.tasks.SourceStreamTask.run(SourceStreamTask.java:94)
~[flink-streaming-java_2.11-1.4.2.jar:1.4.2]
at
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:264)
~[flink-streaming-java_2.11-1.4.2.jar:1.4.2]
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:718)
~[flink-runtime_2.11-1.4.2.jar:1.4.2]
at java.lang.Thread.run(Thread.java:745) [na:1.8.0_60]

when deploy to and run on local cluster -->
2018-03-23 07:27:23.130 [Source: Custom Source -> Map (1/1)] INFO
 org.apache.flink.runtime.taskmanager.Task  - Source: Custom Source -> Map
(1/1) (db21c7604b94968097d4be7b8558ac08) switched from RUNNING to FAILED.
java.lang.RuntimeException: Invalid record:
2264,2013002216,2013002213,START,2013-01-01 00:09:00,1970-01-01
00:00:00,-74.00402,40.742107,-73.98032,40.73522,1
at
com.dataartisans.flinktraining.exercises.datastream_java.datatypes.TaxiRide.fromString(TaxiRide.java:119)
~[blob_p-bc9bbda0acf1e77543b16e023cbf2466bbe42a5d-ac2bf6479faa39e6b27477bb80f66178:na]
at
com.dataartisans.flinktraining.exercises.datastream_java.utils.TaxiRideSchema.deserialize(TaxiRideSchema.java:37)
~[blob_p-bc9bbda0acf1e77543b16e023cbf2466bbe42a5d-ac2bf6479faa39e6b27477bb80f66178:na]
at
com.dataartisans.flinktraining.exercises.datastream_java.utils.TaxiRideSchema.deserialize(TaxiRideSchema.java:28)
~[blob_p-bc9bbda0acf1e77543b16e023cbf2466bbe42a5d-ac2bf6479faa39e6b27477bb80f66178:na]
at
org.apache.flink.streaming.util.serialization.KeyedDeserializationSchemaWrapper.deserialize(KeyedDeserializationSchemaWrapper.java:42)
~[blob_p-bc9bbda0acf1e77543b16e023cbf2466bbe42a5d-ac2bf6479faa39e6b27477bb80f66178:na]
at
org.apache.flink.streaming.connectors.kafka.internal.Kafka09Fetcher.runFetchLoop(Kafka09Fetcher.java:139)
~[blob_p-bc9bbda0acf1e77543b16e023cbf2466bbe42a5d-ac2bf6479faa39e6b27477bb80f66178:na]
at
org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase.run(FlinkKafkaConsumerBase.java:652)
~[blob_p-bc9bbda0acf1e77543b16e023cbf2466bbe42a5d-ac2bf6479faa39e6b27477bb80f66178:na]
at
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:86)
~[flink-dist_2.11-1.4.2.jar:1.4.2]
at
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:55)
~[flink-dist_2.11-1.4.2.jar:1.4.2]
at
org.apache.flink.streaming.runtime.tasks.SourceStreamTask.run(SourceStreamTask.java:94)
~[flink-dist_2.11-1.4.2.jar:1.4.2]
at
org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:264)
~[flink-dist_2.11-1.4.2.jar:1.4.2]
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:718)
~[flink-dist_2.11-1.4.2.jar:1.4.2]
at java.lang.Thread.run(Thread.java:745) [na:1.8.0_60]

I copied the PopularPlacesFromKafka.java from
https://raw.githubusercontent.com/dataArtisans/flink-training-exercises/master/src/main/java/com/dataartisans/flinktraining/exercises/datastream_java/connectors/PopularPlacesFromKafka.java


This is a UTF-8 formatted mail
-----------------------------------------------
James C.-C.Yu
+886988713275

Reply via email to