我使用标准的IO monad.
在某些时候,我需要短路.在给定条件下,我不想运行以下ios.
这是我的解决方案,但我发现它太冗长而且不优雅:
def shortCircuit[A](io: IO[A], continue: Boolean) =
io.map(a => if (continue) Some(a) else None)
for {
a <- io
b <- shortCircuit(io, a == 1)
c <- shortCircuit(io, b.map(_ == 1).getOrElse(false))
d <- shortCircuit(io, b.map(_ == 1).getOrElse(false))
e <- shortCircuit(io, b.map(_ == 1).getOrElse(false))
} yield …
Run Code Online (Sandbox Code Playgroud)
例如,对于第3行,第4行和第5行,我需要重复相同的条件.
有没有更好的办法 ?
我想把一个Try[Option[T]]变成一个Try[T]
这是我的代码
def flattenTry[T](t: Try[Option[T]]) : Try[T] = {
t match {
case f : Failure[T] => f.asInstanceOf[Failure[T]]
case Success(e) =>
e match {
case None => Failure[T](new Exception("Parsing error"))
case Some(s) => Success(s)
}
}
}
Run Code Online (Sandbox Code Playgroud)
有没有更好的办法 ?
给出以下代码:
trait S { type T }
case class A(t: Seq[String]) extends S { type T = Seq[String] }
Run Code Online (Sandbox Code Playgroud)
我不明白这个编译错误:似乎没有使用证据.
def f[S<:A, X](g: => Seq[X])(implicit ev: S#T =:= Seq[X]) = new A(g)
<console>:50: error: type mismatch;
found : Seq[X]
required: Seq[String]
def f[S<:A, X](g: => Seq[X])(implicit ev: S#T =:= Seq[X]) = new A(g)
Run Code Online (Sandbox Code Playgroud) 什么是补充的最佳途径Options到一个List.
这是我的第一次尝试:
def append[A](as: List[A], maybeA1 : Option[A], maybeA2: Option[A]) : List[A] = as ++ maybeA1.toList ++ maybeA2.toList
Run Code Online (Sandbox Code Playgroud)
Con:它创造了2个tmp List
(我知道.toList()是可选的,因为存在从Option [A]到Iterable [A]的隐式转换)
另一种尝试是
def append2[A](ls: List[A], maybeA : Option[A]) : List[A] = maybeA.map(_ :: ls).getOrElse(ls)
def append[A](as: List[A], maybeA1 : Option[A], maybeA2: Option[A]) : List[A] = append2(append2(as, maybeA1), maybeA2)
Run Code Online (Sandbox Code Playgroud)
更好的性能但可读性更低......
还有另外一种方法吗?
该方法的目的是在列表中获取元素,直到达到限制.
例如
我想出了两个不同的实现
def take(l: List[Int], limit: Int): List[Int] = {
var sum = 0
l.takeWhile { e =>
sum += e
sum <= limit
}
}
Run Code Online (Sandbox Code Playgroud)
它很简单,但使用了可变状态.
def take(l: List[Int], limit: Int): List[Int] = {
val summed = l.toStream.scanLeft(0) { case (e, sum) => sum + e }
l.take(summed.indexWhere(_ > limit) - 1)
}
Run Code Online (Sandbox Code Playgroud)
它看起来更干净,但它更冗长,也许内存效率更低,因为需要一个流.
有没有更好的办法 ?
我正在尝试将Spark Streaming 2.2.0与Kafka 0.8一起使用.
我已按照此文档:https: //spark.apache.org/docs/latest/streaming-kafka-0-8-integration.html
但我有一个问题:
[WARN ] 2018-01-25 14:54:01,332 org.apache.spark.scheduler.TaskSetManager - Lost task 3.0 in stage 0.0 (TID 3, ip-10-0-155-42.eu-west-1.compute.internal, executor 8): java.lang.NoSuchMethodError: net.jpountz.util.Utils.checkRange([BII)V
at org.apache.kafka.common.message.KafkaLZ4BlockInputStream.read(KafkaLZ4BlockInputStream.java:176)
Run Code Online (Sandbox Code Playgroud)
关于dependencyGraph,似乎是这样
org.apache.spark:spark-streaming-kafka-0-8_2.11:2.2.0
org.apache.spark:spark-streaming_2.11:2.2.0
org.apache.spark:spark-core_2.11:2.2.0
net.jpountz.lz4:lz4:1.3.0
Run Code Online (Sandbox Code Playgroud)
卡夫卡需要lz4:1.2.0.
[更新]如果我强制lz4的版本为1.2.0.,我有另一个问题
Caused by: java.lang.NoClassDefFoundError: net/jpountz/util/SafeUtils
at org.apache.spark.io.LZ4BlockInputStream.read(LZ4BlockInputStream.java:124)
at java.io.ObjectInputStream$PeekInputStream.read(ObjectInputStream.java:2606)
at java.io.ObjectInputStream$PeekInputStream.readFully(ObjectInputStream.java:2622)
at java.io.ObjectInputStream$BlockDataInputStream.readShort(ObjectInputStream.java:3099)
at java.io.ObjectInputStream.readStreamHeader(ObjectInputStream.java:853)
at java.io.ObjectInputStream.<init>(ObjectInputStream.java:349)
at org.apache.spark.serializer.JavaDeserializationStream$$anon$1.<init>(JavaSerializer.scala:63)
at org.apache.spark.serializer.JavaDeserializationStream.<init>(JavaSerializer.scala:63)
at org.apache.spark.serializer.JavaSerializerInstance.deserializeStream(JavaSerializer.scala:122)
at org.apache.spark.broadcast.TorrentBroadcast$.unBlockifyObject(TorrentBroadcast.scala:291)
at org.apache.spark.broadcast.TorrentBroadcast$$anonfun$readBroadcastBlock$1.apply(TorrentBroadcast.scala:226)
at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1303)
at org.apache.spark.broadcast.TorrentBroadcast.readBroadcastBlock(TorrentBroadcast.scala:206)
at org.apache.spark.broadcast.TorrentBroadcast._value$lzycompute(TorrentBroadcast.scala:66)
at org.apache.spark.broadcast.TorrentBroadcast._value(TorrentBroadcast.scala:66)
at org.apache.spark.broadcast.TorrentBroadcast.getValue(TorrentBroadcast.scala:96)
at org.apache.spark.broadcast.Broadcast.value(Broadcast.scala:70)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:81)
at org.apache.spark.scheduler.Task.run(Task.scala:108)
at …Run Code Online (Sandbox Code Playgroud) 以下是在Spark shell中复制的简单步骤:
scala> case class Foo(d: Option[Double])
defined class Foo
scala> val df = spark.createDataFrame(Seq(Foo(None), Foo(Some(1.0))))
df: org.apache.spark.sql.DataFrame = [d: double]
scala> df.as[Double].printSchema
root
|-- d: double (nullable = true)
scala> df.as[Double].collect
java.lang.NullPointerException: Null value appeared in non-nullable field:
- root class: "scala.Double"
If the schema is inferred from a Scala tuple/case class, or a Java bean, please try to use scala.Option[_] or other nullable types (e.g. java.lang.Integer instead of int/scala.Int).
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificSafeProjection.apply(Unknown Source)
at org.apache.spark.sql.Dataset$$anonfun$org$apache$spark$sql$Dataset$$collectFromPlan$1.apply(Dataset.scala:2864)
at org.apache.spark.sql.Dataset$$anonfun$org$apache$spark$sql$Dataset$$collectFromPlan$1.apply(Dataset.scala:2861)
at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234) …Run Code Online (Sandbox Code Playgroud) 这是我的代码的简化版本。
我怎样才能避免打电话asInstanceOf(因为这是设计不佳的解决方案的味道)?
sealed trait Location
final case class Single(bucket: String) extends Location
final case class Multi(buckets: Seq[String]) extends Location
@SuppressWarnings(Array("org.wartremover.warts.AsInstanceOf"))
class Log[L <: Location](location: L, path: String) { // I prefer composition over inheritance
// I don't want to pass location to this method because it's a property of the object
// It's a separated function because there is another caller
private def getSinglePath()(implicit ev: L <:< Single): String = s"fs://${location.bucket}/$path"
def getPaths(): Seq[String] =
location match { …Run Code Online (Sandbox Code Playgroud) 有没有之间的差异classOf[String].isInstance("42")和"42".isInstanceOf[String]?
如果是的话,你能解释一下吗?
给定一个功能
def f(i: I) : S => S
Run Code Online (Sandbox Code Playgroud)
我想写一个非常常见的组合器 g
def g(is : Seq[I], init: S) : S
Run Code Online (Sandbox Code Playgroud)
简单的实现仅使用经典scala
def g(is : Seq[I], init: S) : S =
is.foldLeft(init){ case (acc, i) => f(i)(acc) }
Run Code Online (Sandbox Code Playgroud)
我尝试使用,Foldable但我遇到了编译问题.
import cats._
import cats.Monoid
import cats.implicits._
def g(is : Seq[I], init: S) : S =
Foldable[List].foldMap(is.toList)(f _)(init)
Run Code Online (Sandbox Code Playgroud)
错误是
could not find implicit value for parameter B: cats.Monoid[S => S]
Run Code Online (Sandbox Code Playgroud)
我成功了 State
import cats.data.State
import cats.instances.all._
import cats.syntax.traverse._
def g(is : Seq[I], init: …Run Code Online (Sandbox Code Playgroud) scala ×8
apache-spark ×2
apache-kafka ×1
casting ×1
collections ×1
implicit ×1
monads ×1
scala-cats ×1
scala-option ×1
typeclass ×1
types ×1