-
Notifications
You must be signed in to change notification settings - Fork 3
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
1 parent
26c7093
commit 0909b6c
Showing
6 changed files
with
131 additions
and
72 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,29 @@ | ||
while true | ||
do | ||
temp=$(shuf -i 18-53 -n 1) | ||
number=$(shuf -i 1-3113 -n 1) | ||
|
||
curl -v -s -S X POST http://localhost:9001 \ | ||
--header 'Content-Type: application/json; charset=utf-8' \ | ||
--header 'Accept: application/json' \ | ||
--header 'User-Agent: orion/0.10.0' \ | ||
--header "Fiware-Service: demo" \ | ||
--header "Fiware-ServicePath: /test" \ | ||
-d '{ | ||
"data": [ | ||
{ | ||
"id": "R1","type": "Node", | ||
"co": {"type": "Float","value": 0,"metadata": {}}, | ||
"co2": {"type": "Float","value": 0,"metadata": {}}, | ||
"humidity": {"type": "Float","value": 40,"metadata": {}}, | ||
"pressure": {"type": "Float","value": '$number',"metadata": {}}, | ||
"temperature": {"type": "Float","value": '$temp',"metadata": {}}, | ||
"wind_speed": {"type": "Float","value": 1.06,"metadata": {}} | ||
} | ||
], | ||
"subscriptionId": "57458eb60962ef754e7c0998" | ||
}' | ||
|
||
|
||
sleep 1 | ||
done |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,50 @@ | ||
while true | ||
do | ||
temp=$(shuf -i 18-53 -n 1) | ||
number=$(shuf -i 1-3113 -n 1) | ||
|
||
curl -v -s -S X POST http://localhost:9001 \ | ||
--header 'Content-Type: application/json; charset=utf-8' \ | ||
--header 'Accept: application/json' \ | ||
--header 'User-Agent: orion/0.10.0' \ | ||
--header "Fiware-Service: demo" \ | ||
--header "Fiware-ServicePath: /test" \ | ||
-d '{ | ||
"data": [ | ||
{ | ||
"id": "R1", | ||
"type": "Node", | ||
"information": { | ||
"type": "object", | ||
"value": { | ||
"buses":[ | ||
{ | ||
"name": "BusCompany1", | ||
"schedule": { | ||
"morning": [7,9,11], | ||
"afternoon": [13,15,17,19], | ||
"night" : [23,1,5] | ||
}, | ||
"price": 19 | ||
}, | ||
{ | ||
"name": "BusCompany2", | ||
"schedule": { | ||
"morning": [8,10,12], | ||
"afternoon": [16,20], | ||
"night" : [23] | ||
}, | ||
"price": 14 | ||
} | ||
] | ||
}, | ||
"metadata": {} | ||
} | ||
} | ||
], | ||
"subscriptionId": "57458eb60962ef754e7c0998" | ||
}' | ||
|
||
|
||
sleep 1 | ||
done |
This file was deleted.
Oops, something went wrong.
9 changes: 4 additions & 5 deletions
9
...ctor/examples/example1/Example1_avg.scala → ...onnector/examples/example4/Example4.scala
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
42 changes: 42 additions & 0 deletions
42
src/main/scala/org/fiware/cosmos/orion/flink/connector/examples/example5/Example5.scala
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,42 @@ | ||
package org.fiware.cosmos.orion.flink.connector.examples.example5 | ||
|
||
import org.apache.flink.api.common.functions.AggregateFunction | ||
import org.apache.flink.streaming.api.scala._ | ||
import org.apache.flink.streaming.api.windowing.time.Time | ||
import org.fiware.cosmos.orion.flink.connector.OrionSource | ||
|
||
/** | ||
* Example5 Orion Connector | ||
* @author @sonsoleslp | ||
*/ | ||
object Example5{ | ||
|
||
def main(args: Array[String]): Unit = { | ||
val env = StreamExecutionEnvironment.getExecutionEnvironment | ||
// Create Orion Source. Receive notifications on port 9001 | ||
val eventStream = env.addSource(new OrionSource(9001)) | ||
|
||
// Process event stream | ||
val processedDataStream = eventStream | ||
.flatMap(event => event.entities) | ||
.map(entity => { | ||
entity.attrs("information").value.asInstanceOf[Map[String, Any]] | ||
}) | ||
.map(list => list("buses").asInstanceOf[List[Map[String,Any]]]) | ||
.flatMap(bus => bus ) | ||
.map(bus => new Bus(bus("name").asInstanceOf[String], | ||
bus("schedule").asInstanceOf[ Map[String, List[ scala.math.BigInt]]], | ||
bus("price").asInstanceOf[ scala.math.BigInt])) | ||
.keyBy("name") | ||
.timeWindow(Time.seconds(5), Time.seconds(2)) | ||
.min("price") | ||
|
||
// print the results with a single thread, rather than in parallel | ||
|
||
processedDataStream.print().setParallelism(1) | ||
|
||
env.execute("Socket Window NgsiEvent") | ||
} | ||
case class Bus(name: String, schedule: Map[String, List[ scala.math.BigInt]], price: scala.math.BigInt) | ||
|
||
} |