【问题标题】:How can we use lagom's Read-side processor with Dgraph?我们如何在 Dgraph 中使用 lagom 的读取端处理器?
【发布时间】:2018-06-03 14:29:17
【问题描述】:

我是 lagom 和 dgraph 的新手。我被困在如何将 lagom 的读取端处理器与 Dgraph 一起使用。下面是使用 Cassandra 和 lagom 的代码。

import akka.NotUsed;
import com.lightbend.lagom.javadsl.api.ServiceCall;
import com.lightbend.lagom.javadsl.persistence.cassandra.CassandraSession;
import java.util.concurrent.CompletableFuture;
import javax.inject.Inject;
import akka.stream.javadsl.Source;
public class FriendServiceImpl implements FriendService {

private final CassandraSession cassandraSession;

@Inject
public FriendServiceImpl(CassandraSession cassandraSession) {
    this.cassandraSession = cassandraSession;
}

//Implement your service method here

}

【问题讨论】:

  • 那么,com.lightbend.lagom.javadsl.persistence.dgraph 是否以包的形式存在?
  • @cricket_007 到目前为止,我还没有遇到您提到的任何此类包。

标签: java microservices lagom dgraph


【解决方案1】:

Lagom 不为 Dgraph 提供开箱即用的支持。如果你必须使用 Lagom 的 Read-Side 处理器和 Dgraph,那么你必须使用 Lagom 的Generic Read Side support。像这样:

/**
 * Read side processor for Dgraph.
 */
public class FriendEventProcessor extends ReadSideProcessor<FriendEvent> {
    private static void createModel() {
        //TODO: Initialize schema in Dgraph
    }

    @Override
    public ReadSideProcessor.ReadSideHandler<FriendEvent> buildHandler() {
        return new ReadSideHandler<FriendEvent>() {
            private final Done doneInstance = Done.getInstance();

            @Override
            public CompletionStage<Done> globalPrepare() {
                createModel();
                return CompletableFuture.completedFuture(doneInstance);
            }

            @Override
            public CompletionStage<Offset> prepare(final AggregateEventTag<FriendEvent> tag) {
                return CompletableFuture.completedFuture(Offset.NONE);
            }

            @Override
            public Flow<Pair<FriendEvent, Offset>, Done, ?> handle() {
                return Flow.<Pair<FriendEvent, Offset>>create()
                        .mapAsync(1, eventAndOffset -> {
                                    if (eventAndOffset.first() instanceof FriendCreated) {
                                        //TODO: Add Friend in Dgraph;
                                    }

                                    return CompletableFuture.completedFuture(doneInstance);
                                }
                        );
            }
        };
    }

    @Override
    public PSequence<AggregateEventTag<FriendEvent>> aggregateTags() {
        return FriendEvent.TAG.allTags();
    }
}

对于FriendEvent.TAG.allTags(),你必须在FriendEvent接口中添加如下代码:

int NUM_SHARDS = 20;

  AggregateEventShards<FriendEvent> TAG =
          AggregateEventTag.sharded(FriendEvent.class, NUM_SHARDS);

  @Override
  default AggregateEventShards<FriendEvent> aggregateTag() {
    return TAG;
  }

我希望这会有所帮助!

【讨论】:

    猜你喜欢
    • 2020-08-30
    • 2020-09-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多