CodisLabs / CodisLabs/codis

Read timed out

Open
#1,413 5 comments 0 reactions 0 assignees View on GitHub
Dominant language
Go
Stars
13.2k
Forks
2.7k
PR merge metrics
No merged PRs in 30d

Description

@yangzhe1991 @elvuel @spinlock @fancy-rabbit 请教个问题,spark streaming消费kafka中的消息,写入codis,执行一段时间后报错,spark streaming程序日志如下:
`User class threw exception: org.apache.spark.SparkException: Job aborted due to stage failure: Task 17 in stage 7813.0 failed 4 times, most recent failure: Lost task 17.3 in stage 7813.0 (TID 468911, DataNode-01, executor 77): redis.clients.jedis.exceptions.JedisConnectionException: java.net.SocketTimeoutException: Read timed out
at redis.clients.util.RedisInputStream.ensureFill(RedisInputStream.java:201)
at redis.clients.util.RedisInputStream.readByte(RedisInputStream.java:40)
at redis.clients.jedis.Protocol.process(Protocol.java:141)
at redis.clients.jedis.Protocol.read(Protocol.java:205)
at redis.clients.jedis.Connection.readProtocolWithCheckingBroken(Connection.java:297)
at redis.clients.jedis.Connection.getAll(Connection.java:267)
at redis.clients.jedis.Connection.getAll(Connection.java:259)
at redis.clients.jedis.Pipeline.sync(Pipeline.java:99)
at com.lqz.sparkStreaming.statics.Statics$$anonfun$doCount$1$$anonfun$apply$1.apply(Statics.scala:131)
at com.lqz.sparkStreaming.statics.Statics$$anonfun$doCount$1$$anonfun$apply$1.apply(Statics.scala:94)
at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$33.apply(RDD.scala:920)
at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$33.apply(RDD.scala:920)
at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1870)
at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1870)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:66)
at org.apache.spark.scheduler.Task.run(Task.scala:89)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:229)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:114`

对应代码处理逻辑片段如下:

static.foreachRDD(rdd` => {
rdd.foreachPartition(partitionOfRecords => {
lazy val pipelineObj = CodisClient.pipelined;
lazy val pipeline = pipelineObj.pipeline;

try {
partitionOfRecords.foreach(record => {
//记录key start
val key = record._1;
pipeline.sadd("test:1:k:" + key.substring(0, key.indexOf(':')), key.substring(key.indexOf(':') + 1));
//记录key end
pipeline.incrBy("test:1:1:" + record._1, record._2._1);
});
} finally {
if (pipelineObj != null) {
pipeline.sync();
pipelineObj.jedis.close();
}
}
})
})

Contributor guide

No contributing guide indexed for this repository

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.