|
@@ -9,7 +9,8 @@ class RouteRfidKafkaAction extends Action {
|
|
|
|
|
|
private function addRfidDataToNingbo( $data, $conn ){
|
|
|
//{"methond":"track","mac":"ffc10063","gps":{"locationState":"A","lat":0,"latType":"N","lng":0,"lngType":"E"},"labels":[{"id":"01000423","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621112931"},{"id":"01000424","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621113202"},{"id":"01000422","event":{"dec":20,"lowBattery":0,"entry":1,"leave":0,"in":1},"time":"160621113303"}]}
|
|
|
- //var_dump($data);
|
|
|
+
|
|
|
+ $this->addMonitorProcess();
|
|
|
|
|
|
if($data['methond']=='login'){
|
|
|
return array('success'=>true,'message'=>'addRfidDataToNingbo failed,methond login!');
|
|
@@ -200,6 +201,7 @@ class RouteRfidKafkaAction extends Action {
|
|
|
$e = oci_error();
|
|
|
trigger_error(htmlentities($e['message'], ENT_QUOTES), E_USER_ERROR);
|
|
|
}
|
|
|
+ echo 11111;
|
|
|
while (true) {
|
|
|
//var_dump($conn);
|
|
|
$message = $consumer->consume(120*1000);
|
|
@@ -614,5 +616,207 @@ class RouteRfidKafkaAction extends Action {
|
|
|
|
|
|
}
|
|
|
|
|
|
+
|
|
|
+ private function addMonitorProcess( ){
|
|
|
+ $redis = Redis("nbfd_monitor_process_id","hash");
|
|
|
+ //$pid=getmypid();
|
|
|
+ $pid=posix_getpid();
|
|
|
+ $key = "monitor_process_id_".$pid;
|
|
|
+ $redisData = array(
|
|
|
+ $key =>json_encode(array(
|
|
|
+ "pid" => $pid,
|
|
|
+ "time" => time(),
|
|
|
+ )
|
|
|
+ )
|
|
|
+ );
|
|
|
+ $redis->add($redisData);
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+ public function pushRfidRouteToStudent( ){
|
|
|
+
|
|
|
+ $broker_list = C('KAFKA_BROKER_LIST');
|
|
|
+
|
|
|
+ if (empty($broker_list)) {
|
|
|
+ exit("KAFKA_BROKER_LIST must be config!".PHP_EOL);
|
|
|
+ }
|
|
|
+ $group = C('ROUTE_INDEX_KAFKA_GROUP_STUDENT');
|
|
|
+ if (empty($group)) {
|
|
|
+ exit("ROUTE_INDEX_KAFKA_GROUP_STUDENT must be config!".PHP_EOL);
|
|
|
+ }
|
|
|
+ $topics = C('ROUTE_INDEX_KAFKA_TOPIC');
|
|
|
+ if (empty($topics)) {
|
|
|
+ exit("ROUTE_INDEX_KAFKA_TOPIC must be config!".PHP_EOL);
|
|
|
+ }
|
|
|
+ $topics = explode(',',$topics);
|
|
|
+
|
|
|
+ // 从 topic :rlstation_rfid_location 取轨迹
|
|
|
+ $conf = new RdKafka\Conf();
|
|
|
+ // Set a rebalance callback to log partition assignments (optional)(当有新的消费进程加入或者退出消费组时,kafka 会自动重新分配分区给消费者进程,这里注册了一个回调函数,当分区被重新分配时触发)
|
|
|
+ $conf->setRebalanceCb(function (RdKafka\KafkaConsumer $kafka, $err, array $partitions = null) {
|
|
|
+ switch ($err) {
|
|
|
+ case RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS:
|
|
|
+ echo "Assign: ";
|
|
|
+ var_dump($partitions);
|
|
|
+ $kafka->assign($partitions);
|
|
|
+ break;
|
|
|
+
|
|
|
+ case RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS:
|
|
|
+ echo "Revoke: ";
|
|
|
+ var_dump($partitions);
|
|
|
+ $kafka->assign(NULL);
|
|
|
+ break;
|
|
|
+
|
|
|
+ default:
|
|
|
+ throw new \Exception($err);
|
|
|
+ }
|
|
|
+ });
|
|
|
+ // Configure the group.id. All consumer with the same group.id will consume( 配置groud.id 具有相同 group.id 的consumer 将会处理不同分区的消息,所以同一个组内的消费者数量如果订阅了一个topic, 那么消费者进程的数量多于这个topic 分区的数量是没有意义的。)
|
|
|
+ // different partitions.
|
|
|
+ $conf->set('group.id', $group);
|
|
|
+ // Initial list of Kafka brokers(添加 kafka集群服务器地址)
|
|
|
+
|
|
|
+ $conf->set('metadata.broker.list', $broker_list);
|
|
|
+ $topicConf = new RdKafka\TopicConf();
|
|
|
+ // Set where to start consuming messages when there is no initial offset in
|
|
|
+ // offset store or the desired offset is out of range.
|
|
|
+ // 'smallest': start from the beginning
|
|
|
+ $topicConf->set('auto.offset.reset', 'latest');
|
|
|
+ // Set the configuration to use for subscribed/assigned topics
|
|
|
+ $conf->setDefaultTopicConf($topicConf);
|
|
|
+ $consumer = new RdKafka\KafkaConsumer($conf);
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+ // 订阅轨迹数据topic
|
|
|
+ $consumer->subscribe($topics);
|
|
|
+
|
|
|
+ $config = C('STUDENT_ORACLE_CONFIG');
|
|
|
+ if (empty($config)) {
|
|
|
+ exit("STUDENT_ORACLE_CONFIG must be config!".PHP_EOL);
|
|
|
+ }
|
|
|
+ $host= $config['host'];
|
|
|
+ $port= $config['port'];
|
|
|
+ $instance_name= $config['instance_name'];
|
|
|
+ $username= $config['username'];
|
|
|
+ $password= $config['password'];
|
|
|
+
|
|
|
+ /*
|
|
|
+ $host= '192.168.100.23';
|
|
|
+ $port= '1521';
|
|
|
+ $instance_name= 'helowin';
|
|
|
+ $username= 'DSSC3';
|
|
|
+ $password= 'Rliandssc3';
|
|
|
+ */
|
|
|
+ $conn = oci_connect($username, $password, $host.':'.$port.'/'. $instance_name,'AL32UTF8');
|
|
|
+ if (!$conn) {
|
|
|
+ $e = oci_error();
|
|
|
+ trigger_error(htmlentities($e['message'], ENT_QUOTES), E_USER_ERROR);
|
|
|
+ }
|
|
|
+ echo 11111;
|
|
|
+ while (true) {
|
|
|
+ //var_dump($conn);
|
|
|
+ $message = $consumer->consume(120*1000);
|
|
|
+
|
|
|
+ switch ($message->err) {
|
|
|
+ case RD_KAFKA_RESP_ERR_NO_ERROR:
|
|
|
+
|
|
|
+ $data = json_decode($message->payload,true);
|
|
|
+ if( $data ){
|
|
|
+
|
|
|
+ $res=$this->addRfidDataToStudent($data,$conn);
|
|
|
+ if(!$res['success']){
|
|
|
+ throw new \Exception($res['message']);
|
|
|
+ }
|
|
|
+ }
|
|
|
+ break;
|
|
|
+ case RD_KAFKA_RESP_ERR__PARTITION_EOF:
|
|
|
+ //echo "No more messages; will wait for more\n";
|
|
|
+ break;
|
|
|
+ case RD_KAFKA_RESP_ERR__TIMED_OUT:
|
|
|
+ echo "Timed out\n";
|
|
|
+ break;
|
|
|
+ default:
|
|
|
+ echo "default break";
|
|
|
+ $this->debug_log( 'default_Log', $message->errstr() );
|
|
|
+ $this->debug_log( 'default_Log', $message->err );
|
|
|
+ break;
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+ private function addRfidDataToStudent( $data, $conn ){
|
|
|
+ //{"methond":"track","mac":"ffc10063","gps":{"locationState":"A","lat":0,"latType":"N","lng":0,"lngType":"E"},"labels":[{"id":"01000423","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621112931"},{"id":"01000424","event":{"dec":8,"lowBattery":0,"entry":0,"leave":1,"in":0},"time":"160621113202"},{"id":"01000422","event":{"dec":20,"lowBattery":0,"entry":1,"leave":0,"in":1},"time":"160621113303"}]}
|
|
|
+
|
|
|
+ if($data['methond']=='login'){
|
|
|
+ return array('success'=>true,'message'=>'addRfidDataToNingbo failed,methond login!');
|
|
|
+ }
|
|
|
+ $RF_ID=strtoupper($data['mac']);
|
|
|
+ $station_cond=array('mac'=>$RF_ID);
|
|
|
+ $device_name=M('stations')->where($station_cond)->getField('name');
|
|
|
+ if(!$device_name){
|
|
|
+ return array('success'=>true,'message'=>'addRfidDataToNingbo failed,station not existed!');
|
|
|
+ }
|
|
|
+ if($data['methond']=='heartbeat'){
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+ if($data['methond']!='track'){
|
|
|
+ return array('success'=>true,'message'=>'addRfidDataToNingbo failed,methond not track!');
|
|
|
+ }
|
|
|
+ if(!$data['labels']){
|
|
|
+ return array('success'=>true,'message'=>'addRfidDataToNingbo failed,labels not existed!');
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+ $updateStationTime=true;
|
|
|
+
|
|
|
+ foreach($data['labels'] as $val){
|
|
|
+ if( ($val['time']<(time()-3600*24*10) ) || ($val['time']>(time()+3600*24*10) ) ){
|
|
|
+ $updateStationTime=false;
|
|
|
+ $val['mac']=$RF_ID;
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ $RF_STAT=0;
|
|
|
+ $plate_no='';
|
|
|
+ if($val['event']['entry']==1){
|
|
|
+ $RF_STAT=1;
|
|
|
+ }elseif($val['event']['leave']==1){
|
|
|
+ $RF_STAT=2;
|
|
|
+ }
|
|
|
+ $RF_FLAGID=strtoupper($val['id']);
|
|
|
+ if($RF_FLAGID=='00000000'){
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ $RF_DATE=date('Y-m-d H:i:s',$val['time']);
|
|
|
+ $sql = 'INSERT INTO "DSSC2"."W_DW_RF_RECORD"("ID", "RF_ID", "RF_FLAGID", "RF_DATE", "RF_STAT") VALUES (DSSC2.SEQ_W_DW_RF_RECORD.nextval, \''.$RF_ID.'\', \''.$RF_FLAGID.'\', TO_DATE(\''.$RF_DATE.'\', \'SYYYY-MM-DD HH24:MI:SS\'), \''.$RF_STAT.'\')';
|
|
|
+ //var_dump($sql);
|
|
|
+ //插入数据到oracle轨迹表
|
|
|
+ $stid = oci_parse($conn, $sql);
|
|
|
+ $r = oci_execute($stid);
|
|
|
+ if(!$r){
|
|
|
+ $val['mac']=$RF_ID;
|
|
|
+ $this->debug_log( 'insert_student_error', $val );
|
|
|
+ return array('success'=>false,'message'=>'addRfidDataToNingbo failed,insert_oracle_error!');
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+ oci_free_statement($stid);
|
|
|
+
|
|
|
+
|
|
|
+ return array('success'=>true,'message'=>'add success');
|
|
|
+ }
|
|
|
+
|
|
|
|
|
|
}
|