【问题标题】:How to subscribe to stream in Angular2?如何在 Angular2 中订阅流?
【发布时间】:2018-02-07 16:28:18
【问题描述】:

我正在 localhost:8080/stream 上运行我的流式 Web 服务,当任何新消息添加到订阅的 mqtt 流时它会响应。我想在我的 Angular2 应用程序中使用这个网络服务。我正在使用 RxJS 在 Angular2 中使用 NodeJS API。我尝试了以下代码,它调用 localhost:8080/stream 一次并结束响应。我希望我的 observable 持续收听网络服务。

var headers = new Headers();
headers.append('Access-Control-Allow-Origin', '*');
headers.append('Content-Type', 'application/json');

let options = new RequestOptions({ headers: headers }); // Create a request option
 return this.http.get("http://localhost:8080/stream", options) // ...using post request
   .map((res: Response) => res.json()) // ...and calling .json() on the response to return data
   .catch((error: any) => Observable.throw(error.json().error));

【问题讨论】:

标签: node.js angular rxjs observable


【解决方案1】:

如果我正确理解您的问题,您希望在某个时间段内使用新消息到达的流中的数据。

要实现这一点,您需要添加订阅服务。

return this.http.get("http://localhost:8080/stream", options) // ...using post request
   .map((res: Response) => res.json()) // ...and calling .json() on the response to return data
   .catch((error: any) => Observable.throw(error.json().error)
   .subscribe(result => this.result =result));

结果会随着新数据的到来而更新,您可以随意使用它。

注意:最佳做法是在服务中单独进行 http 调用并在组件中订阅服务。

为了您的参考,我添加了一个我为演示目的而工作的示例。

  1. 为 http 调用创建服务

@Injectable() 导出类 JsonPlaceHolderService{

    constructor(private http:Http){}

    getAllPosts():Observable<Post[]>{

        return this.http.get("https://jsonplaceholder.typicode.com/posts")
        .map(res=>res.json())


    }       

}

  1. 从您的组件调用服务并不断收听变化。

    导出类 PostsComponent 实现 OnInit{

    constructor(private _service:JsonPlaceHolderService){}
    
    jphPosts:Post[];
    
    title:string="JsonPlaceHolder's Post data";
    
    ngOnInit():void{
    
        this._service.getAllPosts().subscribe(
            (data) =>this.jphPosts = data,
            (err) => console.log(err),
            ()=>console.log("service call completed")
        );
    
    }
    

    }

【讨论】:

  • @VolodymyrBilyachat 如果我遗漏了什么,欢迎您提出建议:)
  • “我希望我的 observable 持续收听网络服务。”所以OP想订阅长池或者websocket
  • 它没有在 websocket 上运行,它只是在 http 服务器上运行。
  • 用 Observable 订阅不就是这样吗?持续监听 api 还是 service?
  • @shaN 当你使用 promise 时,它​​会得到一个响应,但是当使用 Observables 和订阅时,你会在流中得到响应。
【解决方案2】:

您应该在 Angular 上使用 websocket 并让它监听您的服务 URL,然后您应该监听它的事件(打开、关闭、消息),然后使用 Rxjs 创建您自己的主题流以将新数据推送给订阅者。

请查看以下网址:

https://medium.com/@lwojciechowski/websockets-with-angular2-and-rxjs-8b6c5be02fac

【讨论】:

  • 我已经尝试过了,但它对我不起作用,因为它需要 ws url 而我有 http url。
【解决方案3】:

使用 socket.io 将数据从 nodejs 流式传输到 angular

当我尝试这样做时,这会很有用。以下包含来自socket.io package for angular 信用的代码归原作者所有。这取自一个可行的解决方案,可能需要一些调整。

服务器端

var app     = express(),
    http    = require('http'),
    ioServer    = require('socket.io');

var httpServer = http.createServer(app);
var io = new ioServer();
httpServer.listen(1337, function(){
  console.log('httpServer listening on port 1337');
});

io.attach(httpServer);

io.on('connection', function (socket){
    console.log(Connected socket ' + socket.id);
});


//MQTT subscription
client.on('connect', function () { 
    client.subscribe(topic, function () { 
        console.log("subscribed to " + topic) 
        client.on('message', function (topic, msg, pkt) { 

            io.sockets.emit("message", {topic: topic, msg: msg, pkt: pkt});

        }); 
    }); 
});

客户端

使用以下角度创建customService

import * as io from 'socket.io-client';

声明

private ioSocket: any;
private subscribersCounter = 0;

服务类内部

constructor() {

    this.ioSocket = io('socketUrl', {});

}

on(eventName: string, callback: Function) {
    this.ioSocket.on(eventName, callback);
}

removeListener(eventName: string, callback?: Function) {
    return this.ioSocket.removeListener.apply(this.ioSocket, arguments);
}

fromEvent<T>(eventName: string): Observable<T> {
    this.subscribersCounter++;
    return Observable.create((observer: any) => {
        this.ioSocket.on(eventName, (data: T) => {
            observer.next(data);
        });
        return () => {
            if (this.subscribersCounter === 1) {
                this.ioSocket.removeListener(eventName);
            }
        };
    }).share();
}

在您的组件中,将customService 导入为service

service.on("message", dat => {
 console.log(dat);
});

【讨论】:

  • this.client = mqtt.connect() // 你在这里添加一个 ws:// url,但我只有 http url。
  • 两者相同:)
  • 我可以使用我在 NodeJS 中使用的 mqtt 模块
  • 是的,你必须在上面导入 import { * } from 'mqtt';
  • url 会是 mqtt 代理吗?并在 angular2 本身中订阅 mqtt 的主题?
猜你喜欢
  • 2016-03-30
  • 1970-01-01
  • 2017-10-28
  • 2021-08-25
  • 1970-01-01
  • 1970-01-01
  • 2015-09-16
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多