• Flink学习19:算子介绍keyBy


    1.keyBy简介

    主要作用:把相同的数据,汇总到相同的分区中

    (数据本来是分布在不同的slot中,keyBy会把相同的数据拉到相同的slot中)

     

    2.keyBy的使用

    在使用keyBy时候,需要向keyBy传递一个参数,告诉其按照哪个字段进行归类。

    有2种传递参数的方式,

    1.传递位置的数值

    示例:

    import org.apache.flink.api.scala.createTypeInformation
    import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
    
    object keyByTest {
      def main(args: Array[String]): Unit = {
        //create env
        val env = StreamExecutionEnvironment.getExecutionEnvironment
    
        //create ds
        val ds = env.fromElements(("张三", 4), ("张三", 2), ("leo", 5), ("leo", 1),("raj", 8), ("giao", 7))
    
        val keyByedDs = ds.keyBy(0)
    
        keyByedDs.print()
    
        env.execute()
    
    
    
    
    
      }
    
    }
    

    输出结果:

    2.通过名称进行keyBy

    示例:

    import org.apache.flink.api.scala.createTypeInformation
    import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
    
    object keyByNameTest {
    
      //defined the dataSource's data type
      case class StockPrice(stockId:String, timestamp: Long, price:Double)
    
    
      def main(args: Array[String]): Unit = {
    
        //create env
        val env = StreamExecutionEnvironment.getExecutionEnvironment
    
        //create ds
    
        val pricesList = List(StockPrice("stock1", 154545454, 1212.23), StockPrice("stock1", 154545454, 1212.23), StockPrice("stock2", 154545454, 666.23), StockPrice("stock3", 154545454, 888.23))
    
        val ds = env.fromCollection(pricesList)
    
        //transformation
        val keyByedDs = ds.keyBy("stockId")
    
        keyByedDs.print()
    
        env.execute()
    
    
      }
    
    }

    输出结果:

     

  • 相关阅读:
    Anomalib库安装以及使用
    车载以太网物理层SerDes
    学历证书查询 易语言代码
    25期代码随想录算法训练营第一天 | 704. 二分查找,27. 移除元素
    MySQL高级篇05【索引的数据结构】
    MyBatis学习:#占位符和 $占位符的区别
    二分查找算法合集
    ubuntu下载conda
    前端性能优化之数据采集Performance API实战篇
    浅谈设计模式(五)
  • 原文地址:https://blog.csdn.net/hzp666/article/details/126267055