In [1]:
using Rocket
using ReactiveMP
using GraphPPL
using BenchmarkTools
using Distributions
using MacroTools

┌ Info: Precompiling ReactiveMP [a194aa59-28ba-4574-a09c-4a745416d6e3]
└ @ Base loading.jl:1278
│ - If you have ReactiveMP checked out for development and have
│   added Rocket as a dependency but haven't updated your primary
│   environment's manifest file, try `Pkg.resolve()`.
│ - Otherwise you may need to report an issue with ReactiveMP
┌ Info: Precompiling GraphPPL [b3f8163a-e979-4e85-b43e-1f63d8c8b42c]
└ @ Base loading.jl:1278
│ - If you have GraphPPL checked out for development and have
│   added ReactiveMP as a dependency but haven't updated your primary
│   environment's manifest file, try `Pkg.resolve()`.
│ - Otherwise you may need to report an issue with GraphPPL


In [7]:
@model function smoothing(n, k, x0, P)
    
    x_prior ~ NormalMeanVariance(mean(x0), cov(x0)) 

    x = randomvar(n)
    y = datavar(Float64, n)

    x_prev = x_prior

    for i in 1:n
        x[i] ~ (x_prev + 1.0) where { portal = rem(i, k) === 0 ? AsyncPortal() : EmptyPortal() }
        y[i] ~ NormalMeanVariance(x[i], P)
        
        x_prev = x[i]
    end

    return x, y
end

smoothing (generic function with 1 method)

In [8]:
P = 1.0

n = 10_000
k = 500
data = collect(1:n) + rand(Normal(0.0, sqrt(P)), n);

In [9]:
function inference(; data, k, x0, P)
    n = length(data)
    
    _, (x, y) = smoothing(n, k, x0, P);

    buffer    = Vector{Marginal}(undef, n)
    marginals = collectLatest(getmarginals(x))
    
    subscription = subscribe!(marginals, (ms) -> copyto!(buffer, ms))
    
    yield()
    
    update!(y, data)
    
    for i in 1:(div(n, k) + 1)
        yield()
    end
    
    unsubscribe!(subscription)
    
    return buffer
end

inference (generic function with 1 method)

In [10]:
@time res = inference(
    data = data,
    k = k,
    x0 = NormalMeanVariance(0.0, 10000.0),
    P = P
)

  1.450441 seconds (7.54 M allocations: 464.726 MiB, 34.21% gc time)


10000-element Array{Marginal,1}:
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=1.0185512955656066, v=9.999999899999937e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=2.018551295565607, v=9.999999899999937e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=3.018551295565606, v=9.999999899999936e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=4.018551295565607, v=9.999999899999936e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=5.018551295565608, v=9.999999899999937e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=6.018551295565608, v=9.999999899999937e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=7.018551295565609, v=9.999999899999939e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Float64}(μ=8.018551295565608, v=9.999999899999934e-5))
 Marginal{NormalMeanVariance{Float64}}(NormalMeanVariance{Floa