代码之家  ›  专栏  ›  技术社区  ›  Jagrati Gogia

如何使用scala在spark中合并多个数据流?

  •  1
  • Jagrati Gogia  · 技术社区  · 8 年前

    我有三条来自卡夫卡的信息流。我解析作为JSON接收的流,并将其提取到适当的case类中,并形成以下模式的数据流:

    case class Class1(incident_id: String,
                      crt_object_id: String,
                      source: String,
                      order_number: String)
    
    case class Class2(crt_object_id: String,
                      hangup_cause: String)
    
    case class Class3(crt_object_id: String,
                      text: String)
    

    我想基于公共列连接这三个数据流,即。 crt_object_id . 所需数据流的形式应为:

    case class Merged(incident_id: String,
                      crt_object_id: String,
                      source: String,
                      order_number: String,
                      hangup_cause: String,
                      text: String)
    

    请告诉我一个同样的方法。我对Spark和Scala都很陌生。

    1 回复  |  直到 8 年前
        1
  •  2
  •   Alicia Garcia-Raboso    8 年前

    这个 Spark Streaming documentation 告诉您 join 方法:

    join(otherStream, [numTasks])

    2号呼叫时 DStream 第页,共页 (K, V) (K, W) 配对,返回新 数据流 属于 (K, (V, W)) 与每个键的所有元素对配对。

    请注意,您需要 数据流 键值对而非case类的。因此,您必须从case类中提取要加入的字段,加入流并将结果流打包到适当的case类中。

    case class Class1(incident_id: String, crt_object_id: String,
                      source: String, order_number: String)
    case class Class2(crt_object_id: String, hangup_cause: String)
    case class Class3(crt_object_id: String, text: String)
    case class Merged(incident_id: String, crt_object_id: String,
                      source: String, order_number: String,
                      hangup_cause: String, text: String)
    
    val stream1: DStream[Class1] = ...
    val stream2: DStream[Class2] = ...
    val stream3: DStream[Class3] = ...
    
    val transformedStream1: DStream[(String, Class1)] = stream1.map {
        c1 => (c1.crt_object_id, c1)
    }
    val transformedStream2: DStream[(String, Class2)] = stream2.map {
        c2 => (c2.crt_object_id, c2)
    }
    val transformedStream3: DStream[(String, Class3)] = stream3.map {
        c3 => (c3.crt_object_id, c3)
    }
    
    val joined: DStream[(String, ((Class1, Class2), Class3))] =
        transformedStream1.join(transformedStream2).join(transformedStream3)
    
    val merged: DStream[Merged] = joined.map {
        case (crt_object_id, ((c1, c2), c3)) =>
            Merged(c1.incident_id, crt_object_id, c1.source,
                   c1.order_number, c2.hangup_cause, c3.text)
    
    }
    
    推荐文章