1
/*
2
 * Copyright 2017 Spotify AB.
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing,
11
 * software distributed under the License is distributed on an
12
 * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
13
 * KIND, either express or implied.  See the License for the
14
 * specific language governing permissions and limitations
15
 * under the License.
16
 */
17

18
package com.spotify.featran
19

20
import com.esotericsoftware.kryo.serializers.JavaSerializer
21
import org.apache.flink.api.common.typeinfo.TypeInformation
22
import org.apache.flink.api.scala.DataSet
23

24
import scala.reflect.ClassTag
25

26
package object flink {
27

28
  /** [[CollectionType]] for extraction from Apache Flink `DataSet` type. */
29
  implicit object FlinkCollectionType extends CollectionType[DataSet] {
30
    // force fallback to default serializer
31 1
    private val Ti = TypeInformation.of(classOf[Any])
32

33
    override def map[A, B: ClassTag](ma: DataSet[A])(f: A => B): DataSet[B] = {
34 1
      implicit val tib = Ti.asInstanceOf[TypeInformation[B]]
35 1
      ma.map(f)
36
    }
37
    override def reduce[A](ma: DataSet[A])(f: (A, A) => A): DataSet[A] =
38 1
      ma.reduce(f)
39

40
    override def cross[A, B: ClassTag](ma: DataSet[A])(mb: DataSet[B]): DataSet[(A, B)] =
41 1
      ma.crossWithTiny(mb)
42

43
    override def pure[A, B: ClassTag](ma: DataSet[A])(b: B): DataSet[B] = {
44 1
      implicit val tib = Ti.asInstanceOf[TypeInformation[B]]
45 1
      val env = ma.getExecutionEnvironment
46
      // Kryo throws NPE on `Feature`, use Java serialization instead
47 1
      env.addDefaultKryoSerializer(classOf[FeatureSet[Any]], classOf[JavaSerializer])
48 1
      env.fromElements(b)
49
    }
50
  }
51
}

Read our documentation on viewing source code .

Loading