{"id":46868,"date":"2022-11-12T15:46:01","date_gmt":"2023-02-09T02:06:20","guid":{"rendered":"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/"},"modified":"2024-04-29T00:19:24","modified_gmt":"2024-04-28T16:19:24","slug":"%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82","status":"publish","type":"post","link":"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/","title":{"rendered":"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408"},"content":{"rendered":"<p>\u6211\u6253\u7b97\u5c1d\u8bd5\u4f7f\u7528Apache Flink\u8fd9\u4e2a\u6d41\u5904\u7406\u6846\u67b6\uff0c\u7ee7\u7eed\u5b66\u4e60Spark Streaming\u548cKafka Streams\u4e4b\u540e\u3002\u7531\u4e8e\u6211\u5df2\u7ecf\u4f7f\u7528Python\u7f16\u5199\u4e86Spark Streaming\uff0c\u4f7f\u7528Java\u7f16\u5199\u4e86Kafka Streams\uff0c\u6240\u4ee5\u60f3\u5c1d\u8bd5\u7528Scala\u7f16\u5199Apache Flink\u3002<\/p>\n<p>Apache Flink\u4e0eKafka\u4e00\u6837\uff0c\u4e5f\u662f\u4f7f\u7528Scala\u7f16\u5199\u7684\u3002\u5b83\u5f3a\u8c03\u4e86Scala\u7684\u72ec\u7279\u7684\u5411\u540e\u517c\u5bb9\u6027\uff0c\u8fdb\u884c\u7740\u79ef\u6781\u7684\u5f00\u53d1\u3002\u56e0\u6b64\uff0c\u53ef\u4ee5\u5728\u7f51\u7edc\u4e0a\u627e\u5230\u7684\u4fe1\u606f\u5f88\u5feb\u5c31\u4f1a\u8fc7\u65f6\uff0cAPI\u4e5f\u5bb9\u6613\u88ab\u58f0\u660e\u4e3aDeprecated\u6216PublicEvolving\uff0c\u5bf9\u521d\u5b66\u8005\u6765\u8bf4\u6709\u70b9\u96be\u4ee5\u5165\u95e8\u7684\u60c5\u51b5\u3002\u867d\u7136\u5f88\u96be\u627e\u5230\u9002\u5408\u5b66\u4e60\u7684\u597d\u6587\u7ae0\uff0c\u4f46\u6211\u53c2\u8003\u4e86THE RISE OF BIG DATA STREAMING\u4e2d\u7684\u4f20\u611f\u5668\u6570\u636e\u7a97\u53e3\u805a\u5408\u7684\u5199\u6cd5\u3002<\/p>\n<h2>\u9879\u76ee\u6a21\u677f<\/h2>\n<p>\u5728\u4e2d\u56fd\u4eba\u4e2d\u4f7f\u7528Apache Flink\u9879\u76ee\u65f6\uff0c\u4f7f\u7528Scala\u548cSBT\u521b\u5efa\u4e00\u4e2aFlink\u9879\u76ee\u6a21\u677f\u975e\u5e38\u65b9\u4fbf\u3002\u8ba9\u6211\u4eec\u9996\u5148\u4ee5\u6a21\u677f\u4e2d\u7684WordCount\u4f5c\u4e3a\u793a\u4f8b\u6765\u4e86\u89e3\u5982\u4f55\u4f7f\u7528\u3002\u8bf7\u514b\u9686\u8be5\u6a21\u677f\u3002<\/p>\n<pre class=\"post-pre\"><code>$ cd ~\/scala_apps\r\n$ git clone https:\/\/github.com\/tillrohrmann\/flink-project.git\r\n<\/code><\/pre>\n<p>\u8fd9\u91cc\u6709\u51e0\u4e2a\u4f8b\u5b50\uff0c\u4f46\u6211\u4eec\u5c06\u4f7f\u7528WordCount.scala\u3002<\/p>\n<pre class=\"post-pre\"><code>$ tree flink-project\r\nflink-project\/\r\n\u251c\u2500\u2500 build.sbt\r\n\u251c\u2500\u2500 idea.sbt\r\n\u251c\u2500\u2500 project\r\n\u2502   \u251c\u2500\u2500 assembly.sbt\r\n\u2502   \u2514\u2500\u2500 build.properties\r\n\u251c\u2500\u2500 README\r\n\u2514\u2500\u2500 src\r\n    \u2514\u2500\u2500 main\r\n        \u251c\u2500\u2500 resources\r\n        \u2502   \u2514\u2500\u2500 log4j.properties\r\n        \u2514\u2500\u2500 scala\r\n            \u2514\u2500\u2500 org\r\n                \u2514\u2500\u2500 example\r\n                    \u251c\u2500\u2500 Job.scala\r\n                    \u251c\u2500\u2500 SocketTextStreamWordCount.scala\r\n                    \u2514\u2500\u2500 WordCount.scala\r\n\r\n<\/code><\/pre>\n<p>\u5982\u679c\u8981\u4f7f\u7528ENSIME\uff0c\u8bf7\u53c2\u8003\u8fd9\u91cc\u521b\u5efa\u4e00\u4e2a.ensime\u6587\u4ef6\uff0c\u5e76\u5728Emacs\u4e2d\u4f7f\u7528M-x ensime\u547d\u4ee4\u3002<\/p>\n<pre class=\"post-pre\"><code>$ cd ~\/scala_apps\/flink-project\r\n$ sbt\r\n&gt; ensimeConfig\r\n<\/code><\/pre>\n<p>\u4ee5\u4e0b\u662fWordCount.scala\u7684\u4ee3\u7801\u3002\u5b83\u5c06\u8ba1\u7b97\u5305\u542b\u5728\u793a\u4f8b\u6587\u672c\u4e2d\u7684\u82f1\u6587\u5355\u8bcd\u6570\u91cf\u3002<\/p>\n<pre class=\"post-pre\"><code><span class=\"k\">package<\/span> <span class=\"nn\">org.example<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.api.scala._<\/span>\r\n\r\n<span class=\"k\">object<\/span> <span class=\"nc\">WordCount<\/span> <span class=\"o\">{<\/span>\r\n  <span class=\"k\">def<\/span> <span class=\"nf\">main<\/span><span class=\"o\">(<\/span><span class=\"n\">args<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Array<\/span><span class=\"o\">[<\/span><span class=\"kt\">String<\/span><span class=\"o\">])<\/span> <span class=\"o\">{<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">env<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">ExecutionEnvironment<\/span><span class=\"o\">.<\/span><span class=\"py\">getExecutionEnvironment<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">text<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">fromElements<\/span><span class=\"o\">(<\/span><span class=\"s\">\"To be, or not to be,--that is the question:--\"<\/span><span class=\"o\">,<\/span>\r\n      <span class=\"s\">\"Whether 'tis nobler in the mind to suffer\"<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"The slings and arrows of outrageous fortune\"<\/span><span class=\"o\">,<\/span>\r\n      <span class=\"s\">\"Or to take arms against a sea of troubles,\"<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">counts<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">text<\/span><span class=\"o\">.<\/span><span class=\"py\">flatMap<\/span> <span class=\"o\">{<\/span> <span class=\"nv\">_<\/span><span class=\"o\">.<\/span><span class=\"py\">toLowerCase<\/span><span class=\"o\">.<\/span><span class=\"py\">split<\/span><span class=\"o\">(<\/span><span class=\"s\">\"\\\\W+\"<\/span><span class=\"o\">)<\/span> <span class=\"o\">}<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">map<\/span> <span class=\"o\">{<\/span> <span class=\"o\">(<\/span><span class=\"k\">_<\/span><span class=\"o\">,<\/span> <span class=\"mi\">1<\/span><span class=\"o\">)<\/span> <span class=\"o\">}<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">groupBy<\/span><span class=\"o\">(<\/span><span class=\"mi\">0<\/span><span class=\"o\">)<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">sum<\/span><span class=\"o\">(<\/span><span class=\"mi\">1<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"nv\">counts<\/span><span class=\"o\">.<\/span><span class=\"py\">print<\/span><span class=\"o\">()<\/span>\r\n\r\n  <span class=\"o\">}<\/span>\r\n<span class=\"o\">}<\/span>\r\n<\/code><\/pre>\n<p>\u5728\u9879\u76ee\u76ee\u5f55\u4e2d\u8fd0\u884csbt\u7684run\u547d\u4ee4\u3002\u7531\u4e8e\u6709\u51e0\u4e2a\u5b9e\u73b0\u4e86main\u65b9\u6cd5\u7684\u7c7b\uff0c\u6240\u4ee5\u8f93\u5165WordCount 3\u3002<\/p>\n<pre class=\"post-pre\"><code>$ cd ~\/scala_apps\/flink-project\r\n$ sbt\r\n&gt; run\r\nMultiple main classes detected, select one to run:\r\n\r\n [1] org.example.Job\r\n [2] org.example.SocketTextStreamWordCount\r\n [3] org.example.WordCount\r\n\r\nEnter number:3\r\n<\/code><\/pre>\n<p>\u8fd0\u884c\u540e\uff0c\u5c06\u8ba1\u7b97\u5e76\u8f93\u51fa\u6587\u672c\u4e2d\u5305\u542b\u7684\u82f1\u6587\u5355\u8bcd\u6570\u91cf\u3002<\/p>\n<pre class=\"post-pre\"><code> (a,1)\r\n (fortune,1)\r\n (in,1)\r\n (mind,1)\r\n (or,2)\r\n (question,1)\r\n (slings,1)\r\n (suffer,1)\r\n (take,1)\r\n (that,1)\r\n (to,4)\r\n<\/code><\/pre>\n<p>\u6587\u672c\u6570\u636e\u662f\u901a\u8fc7ExecutionEnvironment\u7684fromElements\u65b9\u6cd5\u521b\u5efa\u7684DataSource\u3002<\/p>\n<pre class=\"post-pre\"><code>    <span class=\"k\">val<\/span> <span class=\"nv\">env<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">ExecutionEnvironment<\/span><span class=\"o\">.<\/span><span class=\"py\">getExecutionEnvironment<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">text<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">fromElements<\/span><span class=\"o\">(<\/span><span class=\"s\">\"To be, or not to be,--that is the question:--\"<\/span><span class=\"o\">,<\/span>\r\n      <span class=\"s\">\"Whether 'tis nobler in the mind to suffer\"<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"The slings and arrows of outrageous fortune\"<\/span><span class=\"o\">,<\/span>\r\n      <span class=\"s\">\"Or to take arms against a sea of troubles,\"<\/span><span class=\"o\">)<\/span>\r\n<\/code><\/pre>\n<p>\u5728Apache Flink\u7684Scala\u4e2d\uff0c\u5c3d\u7ba1\u53ef\u4ee5\u66f4\u7b80\u6d01\u5730\u7f16\u5199\u4ee3\u7801\uff0c\u4f46\u6709\u65f6\u5019\u5f88\u96be\u7406\u89e3\u4e0b\u5212\u7ebf\u3001map\u548cgroupBy\u4e2d\u51fa\u73b0\u76840\u548c1\u4ee3\u8868\u7684\u662f\u4ec0\u4e48\u3002\u5728Apache Flink\u7684\u5143\u7ec4\u4e2d\uff0c\u5f53\u901a\u8fc7field\u6307\u5b9a\u65f6\uff0c\u662f\u4ece\u96f6\u5f00\u59cb\u7d22\u5f15\u7684\uff0c\u6240\u4ee5\u987a\u5e8f\u4f9d\u6b21\u662f0\u548c1\u3002<\/p>\n<pre class=\"post-pre\"><code>    <span class=\"k\">val<\/span> <span class=\"nv\">counts<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">text<\/span><span class=\"o\">.<\/span><span class=\"py\">flatMap<\/span> <span class=\"o\">{<\/span> <span class=\"nv\">_<\/span><span class=\"o\">.<\/span><span class=\"py\">toLowerCase<\/span><span class=\"o\">.<\/span><span class=\"py\">split<\/span><span class=\"o\">(<\/span><span class=\"s\">\"\\\\W+\"<\/span><span class=\"o\">)<\/span> <span class=\"o\">}<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">map<\/span> <span class=\"o\">{<\/span> <span class=\"o\">(<\/span><span class=\"k\">_<\/span><span class=\"o\">,<\/span> <span class=\"mi\">1<\/span><span class=\"o\">)<\/span> <span class=\"o\">}<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">groupBy<\/span><span class=\"o\">(<\/span><span class=\"mi\">0<\/span><span class=\"o\">)<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">sum<\/span><span class=\"o\">(<\/span><span class=\"mi\">1<\/span><span class=\"o\">)<\/span>\r\n<\/code><\/pre>\n<p>\u4f7f\u7528flatMap\u548c\u6b63\u5219\u8868\u8fbe\u5f0f\u5c06\u6587\u672c\u6570\u636e\u5206\u5272\u6210\u5355\u8bcd\uff0c\u5e76\u521b\u5efa\u4e00\u4e2aDataSet\u3002\u4f7f\u7528map()\u51fd\u6570\u521b\u5efa\u4e00\u4e2a\u7531\u5355\u8bcd(_)\u548c\u6570\u5b57(1)\u7ec4\u6210\u7684Tuple\u7684DataSet\u3002\u4f7f\u7528groupBy()\u51fd\u6570\uff0c\u5e76\u6307\u5b9a\u5b57\u6bb50\u6765\u6309\u7167\u5355\u8bcd\u8fdb\u884c\u5206\u7ec4\uff0c\u521b\u5efa\u4e00\u4e2a\u6309\u5355\u8bcd\u5206\u7ec4\u7684GroupedDataSet\u3002\u6700\u540e\uff0c\u4f7f\u7528sum()\u51fd\u6570\u7684\u53c2\u6570\u6307\u5b9a\u5b57\u6bb51\uff0c\u5c06\u5355\u8bcd\u548c\u5355\u8bcd\u603b\u6570\u7ec4\u6210\u4e00\u4e2aTuple\u7684AggregateDataSet\u3002<\/p>\n<h2>\u7a97\u53e3\u6c47\u603b<\/h2>\n<p>\u4f7f\u7528\u6b64\u6a21\u677f\u9879\u76ee\uff0c\u7f16\u5199\u4e00\u4e2a\u7a0b\u5e8f\uff0c\u4f7f\u752860\u79d2\u7684\u6eda\u52a8\u7a97\u53e3\u805a\u5408\u4f20\u611f\u5668\u6570\u636e\uff0c\u5e76\u8ba1\u7b97\u5468\u56f4\u6e29\u5ea6(ambient)\u7684\u5e73\u5747\u503c\u3002\u7531\u4e8e\u4f7f\u7528Kafka\u4f5c\u4e3a\u6570\u636e\u6e90\uff0c\u6240\u4ee5\u53c2\u8003\u6b64\u6587\u6863\u5c06\u6765\u81eaRaspberry Pi 3\u7684SensorTag\u6570\u636e\u53d1\u9001\u5230Kafka\u3002<\/p>\n<pre class=\"post-pre\"><code>Raspberry Pi 3 -&gt; Source (Kafka) -&gt; \u30b9\u30c8\u30ea\u30fc\u30e0\u51e6\u7406 -&gt; Sink (Kafka)\r\n<\/code><\/pre>\n<p>\u514b\u9686\u6a21\u677f\u9879\u76ee\u540e\uff0c\u5220\u9664\u73b0\u6709\u6587\u4ef6\u3002<\/p>\n<pre class=\"post-pre\"><code>$ cd ~\/scala_apps\r\n$ git clone https:\/\/github.com\/tillrohrmann\/flink-project.git streams-flink-scala-examples\r\n$ cd streams-flink-scala-examples\r\n$ rm -fr src\/main\/scala\/org\/\r\n<\/code><\/pre>\n<p>\u521b\u5efaScala\u7684\u5305\u76ee\u5f55\u3002<\/p>\n<pre class=\"post-pre\"><code>$ mkdir -p src\/main\/scala\/com\/github\/masato\/streams\/flink\r\n<\/code><\/pre>\n<h3>\u6784\u5efa\u6587\u4ef6.sbt<\/h3>\n<p>Kafka \u4f7f\u7528\u4e0e\u6211\u4eec\u76f8\u540c\u7684 landoop\/fast-data-dev Docker \u6620\u50cf\u3002 \u7248\u672c\u4e3a 0.10.2.1\u3002 \u6211\u4eec\u5c06\u6dfb\u52a0\u652f\u6301 Kafka 0.10 \u7684\u8f6f\u4ef6\u5305\u3002<\/p>\n<pre class=\"post-pre\"><code><span class=\"k\">val<\/span> <span class=\"nv\">flinkVersion<\/span> <span class=\"k\">=<\/span> <span class=\"s\">\"1.3.2\"<\/span>\r\n\r\n<span class=\"k\">val<\/span> <span class=\"nv\">flinkDependencies<\/span> <span class=\"k\">=<\/span> <span class=\"nc\">Seq<\/span><span class=\"o\">(<\/span>\r\n  <span class=\"s\">\"org.apache.flink\"<\/span> <span class=\"o\">%%<\/span> <span class=\"s\">\"flink-scala\"<\/span> <span class=\"o\">%<\/span> <span class=\"n\">flinkVersion<\/span> <span class=\"o\">%<\/span> <span class=\"s\">\"provided\"<\/span><span class=\"o\">,<\/span>\r\n  <span class=\"s\">\"org.apache.flink\"<\/span> <span class=\"o\">%%<\/span> <span class=\"s\">\"flink-streaming-scala\"<\/span> <span class=\"o\">%<\/span> <span class=\"n\">flinkVersion<\/span> <span class=\"o\">%<\/span> <span class=\"s\">\"provided\"<\/span><span class=\"o\">,<\/span>\r\n  <span class=\"s\">\"org.apache.flink\"<\/span> <span class=\"o\">%%<\/span> <span class=\"s\">\"flink-connector-kafka-0.10\"<\/span> <span class=\"o\">%<\/span> <span class=\"n\">flinkVersion<\/span><span class=\"o\">)<\/span>\r\n\r\n<span class=\"k\">val<\/span> <span class=\"nv\">otherDependencies<\/span> <span class=\"k\">=<\/span> <span class=\"nc\">Seq<\/span><span class=\"o\">(<\/span>\r\n   <span class=\"s\">\"com.typesafe\"<\/span> <span class=\"o\">%<\/span> <span class=\"s\">\"config\"<\/span> <span class=\"o\">%<\/span> <span class=\"s\">\"1.3.1\"<\/span>\r\n<span class=\"o\">)<\/span>\r\n\r\n<span class=\"k\">lazy<\/span> <span class=\"k\">val<\/span> <span class=\"nv\">root<\/span> <span class=\"k\">=<\/span> <span class=\"o\">(<\/span><span class=\"n\">project<\/span> <span class=\"n\">in<\/span> <span class=\"nf\">file<\/span><span class=\"o\">(<\/span><span class=\"s\">\".\"<\/span><span class=\"o\">)).<\/span>\r\n  <span class=\"nf\">settings<\/span><span class=\"o\">(<\/span>\r\n    <span class=\"n\">libraryDependencies<\/span> <span class=\"o\">++=<\/span> <span class=\"n\">flinkDependencies<\/span><span class=\"o\">,<\/span>\r\n    <span class=\"n\">libraryDependencies<\/span> <span class=\"o\">++=<\/span> <span class=\"n\">otherDependencies<\/span>\r\n  <span class=\"o\">)<\/span>\r\n\r\n<span class=\"n\">mainClass<\/span> <span class=\"n\">in<\/span> <span class=\"n\">assembly<\/span> <span class=\"o\">:=<\/span> <span class=\"nc\">Some<\/span><span class=\"o\">(<\/span><span class=\"s\">\"com.github.masato.streams.flink.App\"<\/span><span class=\"o\">)<\/span>\r\n<\/code><\/pre>\n<h3>App.scala \u7684\u4e2d\u6587\u91ca\u4e49\u662f\u4ec0\u4e48?<\/h3>\n<p>\u4ee5\u4e0b\u662f\u5b9e\u73b0\u4e86\u4e3b\u8981\u65b9\u6cd5\u7684\u7a0b\u5e8f\u7684\u5b8c\u6574\u5185\u5bb9\u3002Kafka\u8fde\u63a5\u4fe1\u606f\u7b49\u4f7f\u7528config\u8fdb\u884c\u914d\u7f6e\uff0c\u5e76\u5728\u8bbe\u7f6e\u6587\u4ef6\u4e2d\u5b9a\u4e49\u3002\u6e90\u4ee3\u7801\u4e5f\u5b58\u50a8\u5728\u4ee3\u7801\u5e93\u4e2d\u3002<\/p>\n<pre class=\"post-pre\"><code><span class=\"k\">package<\/span> <span class=\"nn\">com.github.masato.streams.flink<\/span>\r\n\r\n<span class=\"k\">import<\/span> <span class=\"nn\">java.util.Properties<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">java.time.ZoneId<\/span><span class=\"o\">;<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">java.time.format.DateTimeFormatter<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">java.util.Date<\/span>\r\n\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.api.TimeCharacteristic<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.api.scala._<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.api.functions.IngestionTimeExtractor<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.api.windowing.time.Time<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.api.windowing.windows.TimeWindow<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.connectors.kafka.<\/span><span class=\"o\">{<\/span><span class=\"nc\">FlinkKafkaConsumer010<\/span><span class=\"o\">,<\/span><span class=\"nc\">FlinkKafkaProducer010<\/span><span class=\"o\">}<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.streaming.util.serialization.<\/span><span class=\"o\">{<\/span><span class=\"nc\">JSONDeserializationSchema<\/span><span class=\"o\">,<\/span><span class=\"nc\">SimpleStringSchema<\/span><span class=\"o\">}<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.api.common.functions.AggregateFunction<\/span>\r\n\r\n<span class=\"k\">import<\/span> <span class=\"nn\">org.apache.flink.util.Collector<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">com.fasterxml.jackson.databind.node.ObjectNode<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">scala.util.parsing.json.JSONObject<\/span>\r\n<span class=\"k\">import<\/span> <span class=\"nn\">com.typesafe.config.ConfigFactory<\/span>\r\n\r\n<span class=\"k\">case<\/span> <span class=\"k\">class<\/span> <span class=\"nc\">Accumulator<\/span><span class=\"o\">(<\/span><span class=\"n\">time<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Long<\/span><span class=\"o\">,<\/span> <span class=\"n\">bid<\/span><span class=\"k\">:<\/span> <span class=\"kt\">String<\/span><span class=\"o\">,<\/span> <span class=\"k\">var<\/span> <span class=\"n\">sum<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Double<\/span><span class=\"o\">,<\/span> <span class=\"k\">var<\/span> <span class=\"n\">count<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Int<\/span><span class=\"o\">)<\/span>\r\n\r\n<span class=\"k\">class<\/span> <span class=\"nc\">Aggregate<\/span> <span class=\"k\">extends<\/span> <span class=\"nc\">AggregateFunction<\/span><span class=\"o\">[(<\/span><span class=\"kt\">String<\/span>, <span class=\"kt\">Double<\/span><span class=\"o\">)<\/span>, <span class=\"kt\">Accumulator<\/span>,<span class=\"kt\">Accumulator<\/span><span class=\"o\">]<\/span> <span class=\"o\">{<\/span>\r\n\r\n  <span class=\"k\">override<\/span> <span class=\"k\">def<\/span> <span class=\"nf\">createAccumulator<\/span><span class=\"o\">()<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Accumulator<\/span> <span class=\"o\">=<\/span> <span class=\"o\">{<\/span>\r\n    <span class=\"k\">return<\/span> <span class=\"nc\">Accumulator<\/span><span class=\"o\">(<\/span><span class=\"mi\">0L<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"\"<\/span><span class=\"o\">,<\/span> <span class=\"mf\">0.0<\/span><span class=\"o\">,<\/span> <span class=\"mi\">0<\/span><span class=\"o\">)<\/span>\r\n  <span class=\"o\">}<\/span>\r\n\r\n  <span class=\"k\">override<\/span> <span class=\"k\">def<\/span> <span class=\"nf\">merge<\/span><span class=\"o\">(<\/span><span class=\"n\">a<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Accumulator<\/span><span class=\"o\">,<\/span> <span class=\"n\">b<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Accumulator<\/span><span class=\"o\">)<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Accumulator<\/span> <span class=\"o\">=<\/span> <span class=\"o\">{<\/span>\r\n    <span class=\"nv\">a<\/span><span class=\"o\">.<\/span><span class=\"py\">sum<\/span> <span class=\"o\">+=<\/span> <span class=\"nv\">b<\/span><span class=\"o\">.<\/span><span class=\"py\">sum<\/span>\r\n    <span class=\"nv\">a<\/span><span class=\"o\">.<\/span><span class=\"py\">count<\/span> <span class=\"o\">+=<\/span> <span class=\"nv\">b<\/span><span class=\"o\">.<\/span><span class=\"py\">count<\/span>\r\n    <span class=\"k\">return<\/span> <span class=\"n\">a<\/span>\r\n  <span class=\"o\">}<\/span>\r\n\r\n  <span class=\"k\">override<\/span> <span class=\"k\">def<\/span> <span class=\"nf\">add<\/span><span class=\"o\">(<\/span><span class=\"n\">value<\/span><span class=\"k\">:<\/span> <span class=\"o\">(<\/span><span class=\"kt\">String<\/span><span class=\"o\">,<\/span> <span class=\"kt\">Double<\/span><span class=\"o\">),<\/span> <span class=\"n\">acc<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Accumulator<\/span><span class=\"o\">)<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Unit<\/span> <span class=\"o\">=<\/span> <span class=\"o\">{<\/span>\r\n    <span class=\"nv\">acc<\/span><span class=\"o\">.<\/span><span class=\"py\">sum<\/span> <span class=\"o\">+=<\/span> <span class=\"nv\">value<\/span><span class=\"o\">.<\/span><span class=\"py\">_2<\/span>\r\n    <span class=\"nv\">acc<\/span><span class=\"o\">.<\/span><span class=\"py\">count<\/span> <span class=\"o\">+=<\/span> <span class=\"mi\">1<\/span>\r\n  <span class=\"o\">}<\/span>\r\n\r\n  <span class=\"k\">override<\/span> <span class=\"k\">def<\/span> <span class=\"nf\">getResult<\/span><span class=\"o\">(<\/span><span class=\"n\">acc<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Accumulator<\/span><span class=\"o\">)<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Accumulator<\/span> <span class=\"o\">=<\/span> <span class=\"o\">{<\/span>\r\n    <span class=\"k\">return<\/span> <span class=\"n\">acc<\/span>\r\n  <span class=\"o\">}<\/span>\r\n<span class=\"o\">}<\/span>\r\n\r\n<span class=\"k\">object<\/span> <span class=\"nc\">App<\/span> <span class=\"o\">{<\/span>\r\n  <span class=\"k\">val<\/span> <span class=\"nv\">fmt<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">DateTimeFormatter<\/span><span class=\"o\">.<\/span><span class=\"py\">ISO_OFFSET_DATE_TIME<\/span>\r\n\r\n  <span class=\"k\">val<\/span> <span class=\"nv\">conf<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">ConfigFactory<\/span><span class=\"o\">.<\/span><span class=\"py\">load<\/span><span class=\"o\">()<\/span>\r\n  <span class=\"k\">val<\/span> <span class=\"nv\">bootstrapServers<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">conf<\/span><span class=\"o\">.<\/span><span class=\"py\">getString<\/span><span class=\"o\">(<\/span><span class=\"s\">\"app.bootstrap-servers\"<\/span><span class=\"o\">)<\/span>\r\n  <span class=\"k\">val<\/span> <span class=\"nv\">groupId<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">conf<\/span><span class=\"o\">.<\/span><span class=\"py\">getString<\/span><span class=\"o\">(<\/span><span class=\"s\">\"app.group-id\"<\/span><span class=\"o\">)<\/span>\r\n  <span class=\"k\">val<\/span> <span class=\"nv\">sourceTopic<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">conf<\/span><span class=\"o\">.<\/span><span class=\"py\">getString<\/span><span class=\"o\">(<\/span><span class=\"s\">\"app.source-topic\"<\/span><span class=\"o\">)<\/span>\r\n  <span class=\"k\">val<\/span> <span class=\"nv\">sinkTopic<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">conf<\/span><span class=\"o\">.<\/span><span class=\"py\">getString<\/span><span class=\"o\">(<\/span><span class=\"s\">\"app.sink-topic\"<\/span><span class=\"o\">)<\/span>\r\n\r\n  <span class=\"k\">def<\/span> <span class=\"nf\">main<\/span><span class=\"o\">(<\/span><span class=\"n\">args<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Array<\/span><span class=\"o\">[<\/span><span class=\"kt\">String<\/span><span class=\"o\">])<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Unit<\/span> <span class=\"o\">=<\/span> <span class=\"o\">{<\/span>\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">props<\/span> <span class=\"k\">=<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">Properties<\/span><span class=\"o\">()<\/span>\r\n    <span class=\"nv\">props<\/span><span class=\"o\">.<\/span><span class=\"py\">setProperty<\/span><span class=\"o\">(<\/span><span class=\"s\">\"bootstrap.servers\"<\/span><span class=\"o\">,<\/span> <span class=\"n\">bootstrapServers<\/span><span class=\"o\">)<\/span>\r\n    <span class=\"nv\">props<\/span><span class=\"o\">.<\/span><span class=\"py\">setProperty<\/span><span class=\"o\">(<\/span><span class=\"s\">\"group.id\"<\/span><span class=\"o\">,<\/span> <span class=\"n\">groupId<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">env<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">StreamExecutionEnvironment<\/span><span class=\"o\">.<\/span><span class=\"py\">getExecutionEnvironment<\/span>\r\n    <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">setStreamTimeCharacteristic<\/span><span class=\"o\">(<\/span><span class=\"nv\">TimeCharacteristic<\/span><span class=\"o\">.<\/span><span class=\"py\">EventTime<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">source<\/span> <span class=\"k\">=<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">FlinkKafkaConsumer010<\/span><span class=\"o\">[<\/span><span class=\"kt\">ObjectNode<\/span><span class=\"o\">](<\/span>\r\n      <span class=\"n\">sourceTopic<\/span><span class=\"o\">,<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">JSONDeserializationSchema<\/span><span class=\"o\">(),<\/span> <span class=\"n\">props<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">events<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">addSource<\/span><span class=\"o\">(<\/span><span class=\"n\">source<\/span><span class=\"o\">).<\/span><span class=\"py\">name<\/span><span class=\"o\">(<\/span><span class=\"s\">\"events\"<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">timestamped<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">events<\/span><span class=\"o\">.<\/span><span class=\"py\">assignTimestampsAndWatermarks<\/span><span class=\"o\">(<\/span>\r\n      <span class=\"k\">new<\/span> <span class=\"nc\">BoundedOutOfOrdernessTimestampExtractor<\/span><span class=\"o\">[<\/span><span class=\"kt\">ObjectNode<\/span><span class=\"o\">](<\/span><span class=\"nv\">Time<\/span><span class=\"o\">.<\/span><span class=\"py\">seconds<\/span><span class=\"o\">(<\/span><span class=\"mi\">10<\/span><span class=\"o\">))<\/span> <span class=\"o\">{<\/span>\r\n        <span class=\"k\">override<\/span> <span class=\"k\">def<\/span> <span class=\"nf\">extractTimestamp<\/span><span class=\"o\">(<\/span><span class=\"n\">element<\/span><span class=\"k\">:<\/span> <span class=\"kt\">ObjectNode<\/span><span class=\"o\">)<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Long<\/span> <span class=\"o\">=<\/span> <span class=\"nv\">element<\/span><span class=\"o\">.<\/span><span class=\"py\">get<\/span><span class=\"o\">(<\/span><span class=\"s\">\"time\"<\/span><span class=\"o\">).<\/span><span class=\"py\">asLong<\/span> <span class=\"o\">*<\/span> <span class=\"mi\">1000<\/span>\r\n      <span class=\"o\">})<\/span>\r\n\r\n    <span class=\"n\">timestamped<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">map<\/span> <span class=\"o\">{<\/span> <span class=\"n\">v<\/span> <span class=\"k\">=&gt;<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">key<\/span> <span class=\"k\">=<\/span>  <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">get<\/span><span class=\"o\">(<\/span><span class=\"s\">\"bid\"<\/span><span class=\"o\">).<\/span><span class=\"py\">asText<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">ambient<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">get<\/span><span class=\"o\">(<\/span><span class=\"s\">\"ambient\"<\/span><span class=\"o\">).<\/span><span class=\"py\">asDouble<\/span>\r\n        <span class=\"o\">(<\/span><span class=\"n\">key<\/span><span class=\"o\">,<\/span> <span class=\"n\">ambient<\/span><span class=\"o\">)<\/span>\r\n      <span class=\"o\">}<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">keyBy<\/span><span class=\"o\">(<\/span><span class=\"n\">v<\/span> <span class=\"k\">=&gt;<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">_1<\/span><span class=\"o\">)<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">timeWindow<\/span><span class=\"o\">(<\/span><span class=\"nv\">Time<\/span><span class=\"o\">.<\/span><span class=\"py\">seconds<\/span><span class=\"o\">(<\/span><span class=\"mi\">60<\/span><span class=\"o\">))<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">aggregate<\/span><span class=\"o\">(<\/span><span class=\"k\">new<\/span> <span class=\"nc\">Aggregate<\/span><span class=\"o\">(),<\/span>\r\n        <span class=\"o\">(<\/span> <span class=\"n\">key<\/span><span class=\"k\">:<\/span> <span class=\"kt\">String<\/span><span class=\"o\">,<\/span>\r\n          <span class=\"n\">window<\/span><span class=\"k\">:<\/span> <span class=\"kt\">TimeWindow<\/span><span class=\"o\">,<\/span>\r\n          <span class=\"n\">input<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Iterable<\/span><span class=\"o\">[<\/span><span class=\"kt\">Accumulator<\/span><span class=\"o\">],<\/span>\r\n          <span class=\"n\">out<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Collector<\/span><span class=\"o\">[<\/span><span class=\"kt\">Accumulator<\/span><span class=\"o\">]<\/span> <span class=\"o\">)<\/span> <span class=\"k\">=&gt;<\/span> <span class=\"o\">{<\/span>\r\n            <span class=\"k\">var<\/span> <span class=\"n\">in<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">input<\/span><span class=\"o\">.<\/span><span class=\"py\">iterator<\/span><span class=\"o\">.<\/span><span class=\"py\">next<\/span><span class=\"o\">()<\/span>\r\n            <span class=\"nv\">out<\/span><span class=\"o\">.<\/span><span class=\"py\">collect<\/span><span class=\"o\">(<\/span><span class=\"nc\">Accumulator<\/span><span class=\"o\">(<\/span><span class=\"nv\">window<\/span><span class=\"o\">.<\/span><span class=\"py\">getEnd<\/span><span class=\"o\">,<\/span> <span class=\"n\">key<\/span><span class=\"o\">,<\/span> <span class=\"nv\">in<\/span><span class=\"o\">.<\/span><span class=\"py\">sum<\/span><span class=\"o\">\/<\/span><span class=\"nv\">in<\/span><span class=\"o\">.<\/span><span class=\"py\">count<\/span><span class=\"o\">,<\/span> <span class=\"nv\">in<\/span><span class=\"o\">.<\/span><span class=\"py\">count<\/span><span class=\"o\">))<\/span>\r\n          <span class=\"o\">}<\/span>\r\n      <span class=\"o\">)<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">map<\/span> <span class=\"o\">{<\/span> <span class=\"n\">v<\/span> <span class=\"k\">=&gt;<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">zdt<\/span> <span class=\"k\">=<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">Date<\/span><span class=\"o\">(<\/span><span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">time<\/span><span class=\"o\">).<\/span><span class=\"py\">toInstant<\/span><span class=\"o\">().<\/span><span class=\"py\">atZone<\/span><span class=\"o\">(<\/span><span class=\"nv\">ZoneId<\/span><span class=\"o\">.<\/span><span class=\"py\">systemDefault<\/span><span class=\"o\">())<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">time<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">fmt<\/span><span class=\"o\">.<\/span><span class=\"py\">format<\/span><span class=\"o\">(<\/span><span class=\"n\">zdt<\/span><span class=\"o\">)<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">json<\/span> <span class=\"k\">=<\/span> <span class=\"nc\">Map<\/span><span class=\"o\">(<\/span><span class=\"s\">\"time\"<\/span> <span class=\"o\">-&gt;<\/span> <span class=\"n\">time<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"bid\"<\/span> <span class=\"o\">-&gt;<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">bid<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"ambient\"<\/span> <span class=\"o\">-&gt;<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">sum<\/span><span class=\"o\">)<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">retval<\/span> <span class=\"k\">=<\/span> <span class=\"nc\">JSONObject<\/span><span class=\"o\">(<\/span><span class=\"n\">json<\/span><span class=\"o\">).<\/span><span class=\"py\">toString<\/span><span class=\"o\">()<\/span>\r\n        <span class=\"nf\">println<\/span><span class=\"o\">(<\/span><span class=\"n\">retval<\/span><span class=\"o\">)<\/span>\r\n        <span class=\"n\">retval<\/span>\r\n      <span class=\"o\">}<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">addSink<\/span><span class=\"o\">(<\/span><span class=\"k\">new<\/span> <span class=\"nc\">FlinkKafkaProducer010<\/span><span class=\"o\">[<\/span><span class=\"kt\">String<\/span><span class=\"o\">](<\/span>\r\n        <span class=\"n\">bootstrapServers<\/span><span class=\"o\">,<\/span>\r\n        <span class=\"n\">sinkTopic<\/span><span class=\"o\">,<\/span>\r\n        <span class=\"k\">new<\/span> <span class=\"nc\">SimpleStringSchema<\/span><span class=\"o\">)<\/span>\r\n      <span class=\"o\">).<\/span><span class=\"py\">name<\/span><span class=\"o\">(<\/span><span class=\"s\">\"kafka\"<\/span><span class=\"o\">)<\/span>\r\n    <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">execute<\/span><span class=\"o\">()<\/span>\r\n  <span class=\"o\">}<\/span>\r\n<span class=\"o\">}<\/span>\r\n<\/code><\/pre>\n<p>\u6211\u4eec\u5c06\u6309\u987a\u5e8f\u67e5\u770bmain()\u51fd\u6570\u7684\u5904\u7406\u3002\u9996\u5148\u8fdb\u884c\u4e0eKafka\u8fde\u63a5\u7684\u8bbe\u7f6e\u3002\u7531\u4e8e\u8fde\u63a5\u7684Kafka\u7248\u672c\u4e3a0.10\uff0c\u6211\u4eec\u5c06\u4f7f\u7528FlinkKafkaConsumer010\u3002\u6765\u81ea\u6811\u8393\u6d3e3\u7684\u4f20\u611f\u5668\u6570\u636e\u91c7\u7528\u4ee5\u4e0bJSON\u683c\u5f0f\u3002<\/p>\n<pre class=\"post-pre\"><code>{'bid': 'B0:B4:48:BE:5E:00', 'time': 1503527847, 'humidity': 26.55792236328125, 'objecttemp': 22.3125, 'ambient': 26.375, 'rh': 76.983642578125}\r\n<\/code><\/pre>\n<p>\u4f7f\u7528JSONDeserializationSchema\u8fdb\u884c\u53cd\u5e8f\u5217\u5316\u3002<\/p>\n<pre class=\"post-pre\"><code>    <span class=\"k\">val<\/span> <span class=\"nv\">props<\/span> <span class=\"k\">=<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">Properties<\/span><span class=\"o\">()<\/span>\r\n    <span class=\"nv\">props<\/span><span class=\"o\">.<\/span><span class=\"py\">setProperty<\/span><span class=\"o\">(<\/span><span class=\"s\">\"bootstrap.servers\"<\/span><span class=\"o\">,<\/span> <span class=\"n\">bootstrapServers<\/span><span class=\"o\">)<\/span>\r\n    <span class=\"nv\">props<\/span><span class=\"o\">.<\/span><span class=\"py\">setProperty<\/span><span class=\"o\">(<\/span><span class=\"s\">\"group.id\"<\/span><span class=\"o\">,<\/span> <span class=\"n\">groupId<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">env<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">StreamExecutionEnvironment<\/span><span class=\"o\">.<\/span><span class=\"py\">getExecutionEnvironment<\/span>\r\n    <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">setStreamTimeCharacteristic<\/span><span class=\"o\">(<\/span><span class=\"nv\">TimeCharacteristic<\/span><span class=\"o\">.<\/span><span class=\"py\">EventTime<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">source<\/span> <span class=\"k\">=<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">FlinkKafkaConsumer010<\/span><span class=\"o\">[<\/span><span class=\"kt\">ObjectNode<\/span><span class=\"o\">](<\/span>\r\n      <span class=\"n\">sourceTopic<\/span><span class=\"o\">,<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">JSONDeserializationSchema<\/span><span class=\"o\">(),<\/span> <span class=\"n\">props<\/span><span class=\"o\">)<\/span>\r\n<\/code><\/pre>\n<p>Apache Flink\u7684\u65f6\u95f4\u6a21\u578b\u88ab\u8bbe\u7f6e\u4e3a\u4e8b\u4ef6\u65f6\u95f4\uff08TimeCharacteristic.EventTime\uff09\u3002\u6211\u4eec\u5c06\u4f20\u611f\u5668\u6570\u636e\u7684time\u5b57\u6bb5\u7528\u4f5c\u65f6\u95f4\u6233\u548c\u6c34\u5370\u3002<\/p>\n<pre class=\"post-pre\"><code>    <span class=\"k\">val<\/span> <span class=\"nv\">events<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">addSource<\/span><span class=\"o\">(<\/span><span class=\"n\">source<\/span><span class=\"o\">).<\/span><span class=\"py\">name<\/span><span class=\"o\">(<\/span><span class=\"s\">\"events\"<\/span><span class=\"o\">)<\/span>\r\n\r\n    <span class=\"k\">val<\/span> <span class=\"nv\">timestamped<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">events<\/span><span class=\"o\">.<\/span><span class=\"py\">assignTimestampsAndWatermarks<\/span><span class=\"o\">(<\/span>\r\n      <span class=\"k\">new<\/span> <span class=\"nc\">BoundedOutOfOrdernessTimestampExtractor<\/span><span class=\"o\">[<\/span><span class=\"kt\">ObjectNode<\/span><span class=\"o\">](<\/span><span class=\"nv\">Time<\/span><span class=\"o\">.<\/span><span class=\"py\">seconds<\/span><span class=\"o\">(<\/span><span class=\"mi\">10<\/span><span class=\"o\">))<\/span> <span class=\"o\">{<\/span>\r\n        <span class=\"k\">override<\/span> <span class=\"k\">def<\/span> <span class=\"nf\">extractTimestamp<\/span><span class=\"o\">(<\/span><span class=\"n\">element<\/span><span class=\"k\">:<\/span> <span class=\"kt\">ObjectNode<\/span><span class=\"o\">)<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Long<\/span> <span class=\"o\">=<\/span> <span class=\"nv\">element<\/span><span class=\"o\">.<\/span><span class=\"py\">get<\/span><span class=\"o\">(<\/span><span class=\"s\">\"time\"<\/span><span class=\"o\">).<\/span><span class=\"py\">asLong<\/span> <span class=\"o\">*<\/span> <span class=\"mi\">1000<\/span>\r\n      <span class=\"o\">})<\/span>\r\n<\/code><\/pre>\n<p>\u4ece\u4f20\u611f\u5668\u4e2d\u53ef\u4ee5\u83b7\u53d6\u5230\u4e00\u4e9b\u6570\u636e\uff0c\u4f46\u5728\u8fd9\u91cc\u6211\u4eec\u53ea\u4f7f\u7528\u73af\u5883\u6e29\u5ea6\u7684\u503c\u3002\u4f7f\u7528map()\u51fd\u6570\u4ee5SensorTag\u7684\u84dd\u7259\u5730\u5740\u4f5c\u4e3a\u952e\u6765\u521b\u5efa\u4e00\u4e2a\u65b0\u7684Tuple\u3002<\/p>\n<pre class=\"post-pre\"><code>    <span class=\"n\">timestamped<\/span>\r\n      <span class=\"o\">.<\/span><span class=\"py\">map<\/span> <span class=\"o\">{<\/span> <span class=\"n\">v<\/span> <span class=\"k\">=&gt;<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">key<\/span> <span class=\"k\">=<\/span>  <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">get<\/span><span class=\"o\">(<\/span><span class=\"s\">\"bid\"<\/span><span class=\"o\">).<\/span><span class=\"py\">asText<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">ambient<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">get<\/span><span class=\"o\">(<\/span><span class=\"s\">\"ambient\"<\/span><span class=\"o\">).<\/span><span class=\"py\">asDouble<\/span>\r\n        <span class=\"o\">(<\/span><span class=\"n\">key<\/span><span class=\"o\">,<\/span> <span class=\"n\">ambient<\/span><span class=\"o\">)<\/span>\r\n      <span class=\"o\">}<\/span>\r\n<\/code><\/pre>\n<p>\u4f7f\u7528DataStream\u7684keyBy()\u65b9\u6cd5\uff0c\u4ee5Tuple\u7684\u7d22\u5f151\u6307\u5b9aBD\u5730\u5740\u4f5c\u4e3a\u952e\u521b\u5efa\u4e86\u4e00\u4e2aKeyedStream\u3002<\/p>\n<pre class=\"post-pre\"><code>      <span class=\"o\">.<\/span><span class=\"py\">keyBy<\/span><span class=\"o\">(<\/span><span class=\"n\">v<\/span> <span class=\"k\">=&gt;<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">_1<\/span><span class=\"o\">)<\/span>\r\n<\/code><\/pre>\n<p>\u4f7f\u7528timeWindow()\u65b9\u6cd5\u521b\u5efa\u4e86\u4e00\u4e2a\u8bbe\u5b9a\u4e3a60\u79d2\u7684\u6eda\u52a8\u7a97\u53e3\u7684WindowedStream\u3002<\/p>\n<pre class=\"post-pre\"><code>      <span class=\"o\">.<\/span><span class=\"py\">timeWindow<\/span><span class=\"o\">(<\/span><span class=\"nv\">Time<\/span><span class=\"o\">.<\/span><span class=\"py\">seconds<\/span><span class=\"o\">(<\/span><span class=\"mi\">60<\/span><span class=\"o\">))<\/span>\r\n<\/code><\/pre>\n<p>\u57281.3\u7248\u672c\u4e2d\uff0capply()\u5df2\u88ab\u5f03\u7528\u3002\u4e4b\u524d\uff0c\u53ef\u4ee5\u8fd9\u6837\u5199\uff1a<\/p>\n<pre class=\"post-pre\"><code>      <span class=\"o\">.<\/span><span class=\"py\">apply<\/span><span class=\"o\">(<\/span>\r\n        <span class=\"o\">(<\/span><span class=\"mi\">0L<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"\"<\/span><span class=\"o\">,<\/span> <span class=\"mf\">0.0<\/span><span class=\"o\">,<\/span> <span class=\"mi\">0<\/span><span class=\"o\">),<\/span>\r\n        <span class=\"o\">(<\/span><span class=\"n\">acc<\/span><span class=\"k\">:<\/span> <span class=\"o\">(<\/span><span class=\"kt\">Long<\/span><span class=\"o\">,<\/span> <span class=\"kt\">String<\/span><span class=\"o\">,<\/span> <span class=\"nc\">Double<\/span><span class=\"o\">,<\/span> <span class=\"nc\">Int<\/span><span class=\"o\">),<\/span>\r\n         <span class=\"n\">v<\/span><span class=\"k\">:<\/span> <span class=\"o\">(<\/span><span class=\"kt\">String<\/span><span class=\"o\">,<\/span> <span class=\"kt\">Double<\/span><span class=\"o\">))<\/span> <span class=\"k\">=&gt;<\/span> <span class=\"o\">{<\/span> <span class=\"o\">(<\/span><span class=\"mi\">0L<\/span><span class=\"o\">,<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">_1<\/span><span class=\"o\">,<\/span> <span class=\"nv\">acc<\/span><span class=\"o\">.<\/span><span class=\"py\">_3<\/span> <span class=\"o\">+<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">_2<\/span><span class=\"o\">,<\/span> <span class=\"nv\">acc<\/span><span class=\"o\">.<\/span><span class=\"py\">_4<\/span> <span class=\"o\">+<\/span> <span class=\"mi\">1<\/span><span class=\"o\">)<\/span> <span class=\"o\">},<\/span>\r\n        <span class=\"o\">(<\/span> <span class=\"n\">window<\/span><span class=\"k\">:<\/span> <span class=\"kt\">TimeWindow<\/span><span class=\"o\">,<\/span>\r\n          <span class=\"n\">counts<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Iterable<\/span><span class=\"o\">[(<\/span><span class=\"kt\">Long<\/span>, <span class=\"kt\">String<\/span>, <span class=\"kt\">Double<\/span>, <span class=\"kt\">Int<\/span><span class=\"o\">)],<\/span>\r\n          <span class=\"n\">out<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Collector<\/span><span class=\"o\">[(<\/span><span class=\"kt\">Long<\/span>, <span class=\"kt\">String<\/span>, <span class=\"kt\">Double<\/span>, <span class=\"kt\">Int<\/span><span class=\"o\">)]<\/span> <span class=\"o\">)<\/span> <span class=\"k\">=&gt;<\/span> <span class=\"o\">{<\/span>\r\n            <span class=\"k\">var<\/span> <span class=\"n\">count<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">counts<\/span><span class=\"o\">.<\/span><span class=\"py\">iterator<\/span><span class=\"o\">.<\/span><span class=\"py\">next<\/span><span class=\"o\">()<\/span>\r\n            <span class=\"nv\">out<\/span><span class=\"o\">.<\/span><span class=\"py\">collect<\/span><span class=\"o\">((<\/span><span class=\"nv\">window<\/span><span class=\"o\">.<\/span><span class=\"py\">getEnd<\/span><span class=\"o\">,<\/span> <span class=\"nv\">count<\/span><span class=\"o\">.<\/span><span class=\"py\">_2<\/span><span class=\"o\">,<\/span> <span class=\"nv\">count<\/span><span class=\"o\">.<\/span><span class=\"py\">_3<\/span><span class=\"o\">\/<\/span><span class=\"nv\">count<\/span><span class=\"o\">.<\/span><span class=\"py\">_4<\/span><span class=\"o\">,<\/span> <span class=\"nv\">count<\/span><span class=\"o\">.<\/span><span class=\"py\">_4<\/span><span class=\"o\">))<\/span>\r\n          <span class=\"o\">}<\/span>\r\n      <span class=\"o\">)<\/span>\r\n<\/code><\/pre>\n<p>\u53e6\u5916\uff0c\u7531\u4e8efold()\u4e5f\u88ab\u6807\u8bb0\u4e3a\u5f03\u7528\uff0c\u6211\u4eec\u63a8\u8350\u4f7f\u7528aggregate()\u8fdb\u884c\u5c1d\u8bd5\u3002aggregate()\u5b9e\u73b0\u4e86AggregateFunction\u3002\u53ef\u4ee5\u50cfapply()\u7684\u793a\u4f8b\u4e00\u6837\u4f7f\u7528\u5143\u7ec4\uff0c\u4f46\u5c06\u5176\u8f6c\u6362\u4e3acase\u7c7b\u4f1a\u66f4\u6613\u8bfb\u4e00\u4e9b\u3002<\/p>\n<pre class=\"post-pre\"><code>      <span class=\"o\">.<\/span><span class=\"py\">aggregate<\/span><span class=\"o\">(<\/span><span class=\"k\">new<\/span> <span class=\"nc\">Aggregate<\/span><span class=\"o\">(),<\/span>\r\n        <span class=\"o\">(<\/span> <span class=\"n\">key<\/span><span class=\"k\">:<\/span> <span class=\"kt\">String<\/span><span class=\"o\">,<\/span>\r\n          <span class=\"n\">window<\/span><span class=\"k\">:<\/span> <span class=\"kt\">TimeWindow<\/span><span class=\"o\">,<\/span>\r\n          <span class=\"n\">input<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Iterable<\/span><span class=\"o\">[<\/span><span class=\"kt\">Accumulator<\/span><span class=\"o\">],<\/span>\r\n          <span class=\"n\">out<\/span><span class=\"k\">:<\/span> <span class=\"kt\">Collector<\/span><span class=\"o\">[<\/span><span class=\"kt\">Accumulator<\/span><span class=\"o\">]<\/span> <span class=\"o\">)<\/span> <span class=\"k\">=&gt;<\/span> <span class=\"o\">{<\/span>\r\n            <span class=\"k\">var<\/span> <span class=\"n\">in<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">input<\/span><span class=\"o\">.<\/span><span class=\"py\">iterator<\/span><span class=\"o\">.<\/span><span class=\"py\">next<\/span><span class=\"o\">()<\/span>\r\n            <span class=\"nv\">out<\/span><span class=\"o\">.<\/span><span class=\"py\">collect<\/span><span class=\"o\">(<\/span><span class=\"nc\">Accumulator<\/span><span class=\"o\">(<\/span><span class=\"nv\">window<\/span><span class=\"o\">.<\/span><span class=\"py\">getEnd<\/span><span class=\"o\">,<\/span> <span class=\"n\">key<\/span><span class=\"o\">,<\/span> <span class=\"nv\">in<\/span><span class=\"o\">.<\/span><span class=\"py\">sum<\/span><span class=\"o\">\/<\/span><span class=\"nv\">in<\/span><span class=\"o\">.<\/span><span class=\"py\">count<\/span><span class=\"o\">,<\/span> <span class=\"nv\">in<\/span><span class=\"o\">.<\/span><span class=\"py\">count<\/span><span class=\"o\">))<\/span>\r\n          <span class=\"o\">}<\/span>\r\n      <span class=\"o\">)<\/span>\r\n<\/code><\/pre>\n<p>\u4e3a\u4e86\u65b9\u4fbf\u4e0e\u5916\u90e8\u7cfb\u7edf\u8fdb\u884c\u534f\u4f5c\uff0c\u6570\u636e\u6d41\u4f1a\u88ab\u6620\u5c04\u4e3a\u5e26\u6709\u65f6\u533a\u7684ISO-8601\u683c\u5f0f\u7684JSON\u5b57\u7b26\u4e32\uff0c\u5e76\u4f7f\u7528UNIX\u65f6\u95f4\u8fdb\u884cmap()\u64cd\u4f5c\u3002\u8fd9\u91cc\u4e3a\u4e86\u8c03\u8bd5\u76ee\u7684\uff0cJSON\u5b57\u7b26\u4e32\u88ab\u8f93\u51fa\u5230\u6807\u51c6\u8f93\u51fa\u3002<\/p>\n<pre class=\"post-pre\"><code>      <span class=\"o\">.<\/span><span class=\"py\">map<\/span> <span class=\"o\">{<\/span> <span class=\"n\">v<\/span> <span class=\"k\">=&gt;<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">zdt<\/span> <span class=\"k\">=<\/span> <span class=\"k\">new<\/span> <span class=\"nc\">Date<\/span><span class=\"o\">(<\/span><span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">time<\/span><span class=\"o\">).<\/span><span class=\"py\">toInstant<\/span><span class=\"o\">().<\/span><span class=\"py\">atZone<\/span><span class=\"o\">(<\/span><span class=\"nv\">ZoneId<\/span><span class=\"o\">.<\/span><span class=\"py\">systemDefault<\/span><span class=\"o\">())<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">time<\/span> <span class=\"k\">=<\/span> <span class=\"nv\">fmt<\/span><span class=\"o\">.<\/span><span class=\"py\">format<\/span><span class=\"o\">(<\/span><span class=\"n\">zdt<\/span><span class=\"o\">)<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">json<\/span> <span class=\"k\">=<\/span> <span class=\"nc\">Map<\/span><span class=\"o\">(<\/span><span class=\"s\">\"time\"<\/span> <span class=\"o\">-&gt;<\/span> <span class=\"n\">time<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"bid\"<\/span> <span class=\"o\">-&gt;<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">bid<\/span><span class=\"o\">,<\/span> <span class=\"s\">\"ambient\"<\/span> <span class=\"o\">-&gt;<\/span> <span class=\"nv\">v<\/span><span class=\"o\">.<\/span><span class=\"py\">sum<\/span><span class=\"o\">)<\/span>\r\n        <span class=\"k\">val<\/span> <span class=\"nv\">retval<\/span> <span class=\"k\">=<\/span> <span class=\"nc\">JSONObject<\/span><span class=\"o\">(<\/span><span class=\"n\">json<\/span><span class=\"o\">).<\/span><span class=\"py\">toString<\/span><span class=\"o\">()<\/span>\r\n        <span class=\"nf\">println<\/span><span class=\"o\">(<\/span><span class=\"n\">retval<\/span><span class=\"o\">)<\/span>\r\n        <span class=\"n\">retval<\/span>\r\n      <span class=\"o\">}<\/span>\r\n<\/code><\/pre>\n<p>\u5728\u6700\u540e\uff0c\u5c06\u6570\u636e\u6d41Sink\u5230Kafka\u4e2d\u3002\u5982\u679c\u7ed9Sink\u547d\u540d\u4e3aname\uff08&#8221;kafka&#8221;\uff09\uff0c\u5219\u4f1a\u5728\u8fd0\u884c\u65f6\u7684\u65e5\u5fd7\u4e2d\u663e\u793a\u3002<\/p>\n<pre class=\"post-pre\"><code>      <span class=\"o\">.<\/span><span class=\"py\">addSink<\/span><span class=\"o\">(<\/span><span class=\"k\">new<\/span> <span class=\"nc\">FlinkKafkaProducer010<\/span><span class=\"o\">[<\/span><span class=\"kt\">String<\/span><span class=\"o\">](<\/span>\r\n        <span class=\"n\">bootstrapServers<\/span><span class=\"o\">,<\/span>\r\n        <span class=\"n\">sinkTopic<\/span><span class=\"o\">,<\/span>\r\n        <span class=\"k\">new<\/span> <span class=\"nc\">SimpleStringSchema<\/span><span class=\"o\">)<\/span>\r\n      <span class=\"o\">).<\/span><span class=\"py\">name<\/span><span class=\"o\">(<\/span><span class=\"s\">\"kafka\"<\/span><span class=\"o\">)<\/span>\r\n    <span class=\"nv\">env<\/span><span class=\"o\">.<\/span><span class=\"py\">execute<\/span><span class=\"o\">()<\/span>\r\n<\/code><\/pre>\n<h3>\u8fd0\u884csbt<\/h3>\n<p>\u6211\u5c06\u8f6c\u5230\u9879\u76ee\u5e76\u6267\u884csbt\u7684run\u547d\u4ee4\u3002<\/p>\n<pre class=\"post-pre\"><code>$ cd ~\/scala_apps\/streams-flink-scala-examples\r\n$ sbt\r\n&gt; run\r\n<\/code><\/pre>\n<p>\u901a\u8fc760\u79d2\u7684\u6eda\u52a8\u7a97\u53e3\u5bf9\u5468\u56f4\u73af\u5883\u6e29\u5ea6\u8fdb\u884c\u4e86\u5e73\u5747\u503c\u7684\u6536\u96c6\uff0c\u5e76\u5c06\u7ed3\u679c\u8f93\u51fa\u5230\u6807\u51c6\u8f93\u51fa\u3002<\/p>\n<pre class=\"post-pre\"><code>{\"time\" : \"2017-08-24T08:10:00+09:00\", \"bid\" : \"B0:B4:48:BE:5E:00\", \"ambient\" : 26.203125}\r\n{\"time\" : \"2017-08-24T08:11:00+09:00\", \"bid\" : \"B0:B4:48:BE:5E:00\", \"ambient\" : 26.234375}\r\n{\"time\" : \"2017-08-24T08:12:00+09:00\", \"bid\" : \"B0:B4:48:BE:5E:00\", \"ambient\" : 26.26875}\r\n<\/code><\/pre>\n","protected":false},"excerpt":{"rendered":"<p>\u6211\u6253\u7b97\u5c1d\u8bd5\u4f7f\u7528Apache Flink\u8fd9\u4e2a\u6d41\u5904\u7406\u6846\u67b6\uff0c\u7ee7\u7eed\u5b66\u4e60Spark Streaming\u548cKafka St [&hellip;]<\/p>\n","protected":false},"author":9,"featured_media":0,"comment_status":"closed","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[1],"tags":[],"class_list":["post-46868","post","type-post","status-publish","format-standard","hentry","category-uncategorized"],"yoast_head":"<!-- This site is optimized with the Yoast SEO Premium plugin v21.5 (Yoast SEO v21.5) - https:\/\/yoast.com\/wordpress\/plugins\/seo\/ -->\n<title>\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408 - Blog - Silicon Cloud<\/title>\n<meta name=\"robots\" content=\"index, follow, max-snippet:-1, max-image-preview:large, max-video-preview:-1\" \/>\n<link rel=\"canonical\" href=\"https:\/\/www.silicloud.com\/zh\/blog\/\u4f7f\u7528apache-flink\u548cscala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408\u3002\/\" \/>\n<meta property=\"og:locale\" content=\"zh_CN\" \/>\n<meta property=\"og:type\" content=\"article\" \/>\n<meta property=\"og:title\" content=\"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408\" \/>\n<meta property=\"og:description\" content=\"\u6211\u6253\u7b97\u5c1d\u8bd5\u4f7f\u7528Apache Flink\u8fd9\u4e2a\u6d41\u5904\u7406\u6846\u67b6\uff0c\u7ee7\u7eed\u5b66\u4e60Spark Streaming\u548cKafka St [&hellip;]\" \/>\n<meta property=\"og:url\" content=\"https:\/\/www.silicloud.com\/zh\/blog\/\u4f7f\u7528apache-flink\u548cscala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408\u3002\/\" \/>\n<meta property=\"og:site_name\" content=\"Blog - Silicon Cloud\" \/>\n<meta property=\"article:published_time\" content=\"2023-02-09T02:06:20+00:00\" \/>\n<meta property=\"article:modified_time\" content=\"2024-04-28T16:19:24+00:00\" \/>\n<meta name=\"author\" content=\"\u6e05, \u626c\" \/>\n<meta name=\"twitter:card\" content=\"summary_large_image\" \/>\n<meta name=\"twitter:label1\" content=\"\u4f5c\u8005\" \/>\n\t<meta name=\"twitter:data1\" content=\"\u6e05, \u626c\" \/>\n\t<meta name=\"twitter:label2\" content=\"\u9884\u8ba1\u9605\u8bfb\u65f6\u95f4\" \/>\n\t<meta name=\"twitter:data2\" content=\"5 \u5206\" \/>\n<script type=\"application\/ld+json\" class=\"yoast-schema-graph\">{\"@context\":\"https:\/\/schema.org\",\"@graph\":[{\"@type\":\"WebPage\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/\",\"url\":\"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/\",\"name\":\"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408 - Blog - Silicon Cloud\",\"isPartOf\":{\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#website\"},\"datePublished\":\"2023-02-09T02:06:20+00:00\",\"dateModified\":\"2024-04-28T16:19:24+00:00\",\"author\":{\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461\"},\"breadcrumb\":{\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/#breadcrumb\"},\"inLanguage\":\"zh-Hans\",\"potentialAction\":[{\"@type\":\"ReadAction\",\"target\":[\"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/\"]}]},{\"@type\":\"BreadcrumbList\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/#breadcrumb\",\"itemListElement\":[{\"@type\":\"ListItem\",\"position\":1,\"name\":\"\u9996\u9875\",\"item\":\"https:\/\/www.silicloud.com\/zh\/blog\/\"},{\"@type\":\"ListItem\",\"position\":2,\"name\":\"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408\"}]},{\"@type\":\"WebSite\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#website\",\"url\":\"https:\/\/www.silicloud.com\/zh\/blog\/\",\"name\":\"Blog - Silicon Cloud\",\"description\":\"\",\"inLanguage\":\"zh-Hans\"},{\"@type\":\"Person\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461\",\"name\":\"\u6e05, \u626c\",\"image\":{\"@type\":\"ImageObject\",\"inLanguage\":\"zh-Hans\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/image\/\",\"url\":\"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g\",\"contentUrl\":\"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g\",\"caption\":\"\u6e05, \u626c\"},\"url\":\"https:\/\/www.silicloud.com\/zh\/blog\/author\/qingyang\/\"},{\"@type\":\"ImageObject\",\"inLanguage\":\"zh-Hans\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/#local-main-organization-logo\",\"url\":\"\",\"contentUrl\":\"\",\"caption\":\"Blog - Silicon Cloud\"}]}<\/script>\n<!-- \/ Yoast SEO Premium plugin. -->","yoast_head_json":{"title":"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408 - Blog - Silicon Cloud","robots":{"index":"index","follow":"follow","max-snippet":"max-snippet:-1","max-image-preview":"max-image-preview:large","max-video-preview":"max-video-preview:-1"},"canonical":"https:\/\/www.silicloud.com\/zh\/blog\/\u4f7f\u7528apache-flink\u548cscala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408\u3002\/","og_locale":"zh_CN","og_type":"article","og_title":"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408","og_description":"\u6211\u6253\u7b97\u5c1d\u8bd5\u4f7f\u7528Apache Flink\u8fd9\u4e2a\u6d41\u5904\u7406\u6846\u67b6\uff0c\u7ee7\u7eed\u5b66\u4e60Spark Streaming\u548cKafka St [&hellip;]","og_url":"https:\/\/www.silicloud.com\/zh\/blog\/\u4f7f\u7528apache-flink\u548cscala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408\u3002\/","og_site_name":"Blog - Silicon Cloud","article_published_time":"2023-02-09T02:06:20+00:00","article_modified_time":"2024-04-28T16:19:24+00:00","author":"\u6e05, \u626c","twitter_card":"summary_large_image","twitter_misc":{"\u4f5c\u8005":"\u6e05, \u626c","\u9884\u8ba1\u9605\u8bfb\u65f6\u95f4":"5 \u5206"},"schema":{"@context":"https:\/\/schema.org","@graph":[{"@type":"WebPage","@id":"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/","url":"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/","name":"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408 - Blog - Silicon Cloud","isPartOf":{"@id":"https:\/\/www.silicloud.com\/zh\/blog\/#website"},"datePublished":"2023-02-09T02:06:20+00:00","dateModified":"2024-04-28T16:19:24+00:00","author":{"@id":"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461"},"breadcrumb":{"@id":"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/#breadcrumb"},"inLanguage":"zh-Hans","potentialAction":[{"@type":"ReadAction","target":["https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/"]}]},{"@type":"BreadcrumbList","@id":"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/#breadcrumb","itemListElement":[{"@type":"ListItem","position":1,"name":"\u9996\u9875","item":"https:\/\/www.silicloud.com\/zh\/blog\/"},{"@type":"ListItem","position":2,"name":"\u4f7f\u7528Apache Flink\u548cScala\u5bf9\u4f20\u611f\u5668\u6570\u636e\u8fdb\u884c\u7a97\u53e3\u805a\u5408"}]},{"@type":"WebSite","@id":"https:\/\/www.silicloud.com\/zh\/blog\/#website","url":"https:\/\/www.silicloud.com\/zh\/blog\/","name":"Blog - Silicon Cloud","description":"","inLanguage":"zh-Hans"},{"@type":"Person","@id":"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461","name":"\u6e05, \u626c","image":{"@type":"ImageObject","inLanguage":"zh-Hans","@id":"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/image\/","url":"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g","contentUrl":"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g","caption":"\u6e05, \u626c"},"url":"https:\/\/www.silicloud.com\/zh\/blog\/author\/qingyang\/"},{"@type":"ImageObject","inLanguage":"zh-Hans","@id":"https:\/\/www.silicloud.com\/zh\/blog\/%e4%bd%bf%e7%94%a8apache-flink%e5%92%8cscala%e5%af%b9%e4%bc%a0%e6%84%9f%e5%99%a8%e6%95%b0%e6%8d%ae%e8%bf%9b%e8%a1%8c%e7%aa%97%e5%8f%a3%e8%81%9a%e5%90%88%e3%80%82\/#local-main-organization-logo","url":"","contentUrl":"","caption":"Blog - Silicon Cloud"}]}},"_links":{"self":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts\/46868","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/users\/9"}],"replies":[{"embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/comments?post=46868"}],"version-history":[{"count":2,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts\/46868\/revisions"}],"predecessor-version":[{"id":67814,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts\/46868\/revisions\/67814"}],"wp:attachment":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/media?parent=46868"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/categories?post=46868"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/tags?post=46868"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}