• 设为首页
  • 点击收藏
  • 手机版
    手机扫一扫访问
    迪恩网络手机版
  • 关注官方公众号
    微信扫一扫关注
    迪恩网络公众号

Java Buffer类代码示例

原作者: [db:作者] 来自: [db:来源] 收藏 邀请

本文整理汇总了Java中io.vertx.rxjava.core.buffer.Buffer的典型用法代码示例。如果您正苦于以下问题:Java Buffer类的具体用法?Java Buffer怎么用?Java Buffer使用的例子?那么恭喜您, 这里精选的类代码示例或许可以为您提供帮助。



Buffer类属于io.vertx.rxjava.core.buffer包,在下文中一共展示了Buffer类的20个代码示例,这些例子默认根据受欢迎程度排序。您可以为喜欢或者感觉有用的代码点赞,您的评价将有助于我们的系统推荐出更棒的Java代码示例。

示例1: getShoppingList

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
private void getShoppingList(RoutingContext rc) {
    Single<HttpResponse<Buffer>> list = shopping.get("/shopping").rxSend();

    HttpServerResponse serverResponse = rc.response()
        .setChunked(true);

    list.subscribe(
        response -> {
            JsonObject body = response.bodyAsJsonObject();
            Observable.from(body)
                .flatMap(entry ->
                    retrievePrice(pricer, entry).toObservable())
                .subscribe(
                    product -> writeProductLine(serverResponse, product),
                    rc::fail,
                    serverResponse::end
                );
        },
        rc::fail);
}
 
开发者ID:cescoffier,项目名称:vertx-chtijug-2017,代码行数:21,代码来源:MyShoppingList.java


示例2: getShoppingListWithCB

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
private void getShoppingListWithCB(RoutingContext rc) {
    Single<HttpResponse<Buffer>> list = shopping.get("/shopping").rxSend();

    HttpServerResponse serverResponse = rc.response()
        .setChunked(true);

    list.subscribe(
        response -> {
            JsonObject body = response.bodyAsJsonObject();
            Observable.from(body)
                .flatMap(entry ->
                    circuit.executeWithFallback(
                        future -> retrievePrice(pricer, entry, future),
                        t -> getFallbackPrice(entry)
                    )
                        .rxSetHandler().toObservable())
                .subscribe(
                    product -> writeProductLine(serverResponse, product),
                    rc::fail,
                    serverResponse::end
                );
        },
        rc::fail);
}
 
开发者ID:cescoffier,项目名称:vertx-chtijug-2017,代码行数:25,代码来源:MyShoppingList.java


示例3: testGetHistory

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Test
public void testGetHistory(TestContext ctx) throws Exception {
    localDatabase.updateOrInsertBlocking(SQL_INSERT_MOVIES_HOBBIT, null);
    localDatabase.updateOrInsertBlocking(SQL_INSERT_MOVIES_GHOST, null);
    localDatabase.updateOrInsertBlocking(SQL_INSERT_VIEW_HOBBIT, null);
    localDatabase.updateOrInsertBlocking(SQL_INSERT_VIEW_GHOST, null);
    Async async = ctx.async();
    doRequest(client.get(API_HISTORY), ctx, cookies, new JsonObject()
            .put("is-first", false)
            .put("is-cinema", false)
            .put("start", "1 January, 2017")
            .put("end", "31 December, 2017")
            .put("page", 0).encode())
            .doOnNext(res -> assertThat(res.statusCode(), is(Status.OK)))
            .flatMap(Utils::bodyToObservable)
            .map(Buffer::toJsonObject)
            .map(DatabaseService::getRows)
            .doOnNext(results -> assertThat(results.size(), is(2)))
            .doOnNext(results -> results.stream()
                    .map(JsonObj::fromParent)
                    .forEach(view -> check(view.getInteger("Id") == 1,
                            () -> assertHobbitView(view),
                            () -> assertGhostView(view))))
            .subscribe(results -> async.complete());
}
 
开发者ID:kristenkotkas,项目名称:moviediary,代码行数:26,代码来源:DatabaseRouterTest.java


示例4: observeGetObject

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
private Observable<ObjectNode> observeGetObject(String endpoint) {
    ObservableHandler<HttpClientResponse> responseObservable = RxHelper.observableHandler();
    client.getNow(endpoint, responseObservable.toHandler());
    return responseObservable
        .flatMap(response -> {
            if (response.statusCode() == HTTP_OK) {
                ObservableHandler<Buffer> observable = RxHelper.observableHandler();
                response.bodyHandler(observable.toHandler());
                return observable;
            } else {
                throw new RuntimeException(String.format(
                        "Endpoint '%s' get returned %d status",
                        endpoint, response.statusCode()));
            }
        })
        .map(buffer -> Json.decodeValue(buffer.toString(UTF_8.name()), ObjectNode.class));
}
 
开发者ID:MiguelGL,项目名称:popular-purchases-demo,代码行数:18,代码来源:RecentPurchasesHandler.java


示例5: importJSON

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
/**
 * Imports a JSON file from the given input stream into the store
 * @param f the JSON file to read
 * @param correlationId a unique identifier for this import process
 * @param filename the name of the file currently being imported
 * @param timestamp denotes when the import process has started
 * @param layer the layer where the file should be stored (may be null)
 * @param tags the list of tags to attach to the file (may be null)
 * @param properties the map of properties to attach to the file (may be null)
 * @return a single that will emit when the file has been imported
 */
protected Single<Integer> importJSON(ReadStream<Buffer> f, String correlationId,
    String filename, long timestamp, String layer, List<String> tags, Map<String, Object> properties) {
  UTF8BomFilter bomFilter = new UTF8BomFilter();
  StringWindow window = new StringWindow();
  GeoJsonSplitter splitter = new GeoJsonSplitter(window);
  AtomicInteger processing = new AtomicInteger(0);
  return f.toObservable()
      .map(buf -> (io.vertx.core.buffer.Buffer)buf.getDelegate())
      .map(bomFilter::filter)
      .doOnNext(window::append)
      .lift(new JsonParserOperator())
      .flatMap(splitter::onEventObservable)
      .flatMapSingle(result -> {
        IndexMeta indexMeta = new IndexMeta(correlationId, filename,
            timestamp, tags, properties, null);
        return addToStoreWithPause(result, layer, indexMeta, f, processing);
      })
      .count()
      .toSingle();
}
 
开发者ID:georocket,项目名称:georocket,代码行数:32,代码来源:ImporterVerticle.java


示例6: addToStoreWithPause

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
/**
 * Add a chunk to the store. Pause the given read stream before adding and
 * increase the given counter. Decrease the counter after the chunk has been
 * written and only resume the read stream if the counter is <code>0</code>.
 * This is necessary because the writing to the store may take longer than
 * reading. We need to pause reading so the store is not overloaded (i.e.
 * we handle back-pressure here).
 * @param chunk the chunk to write
 * @param layer the layer the chunk should be added to (may be null)
 * @param indexMeta metadata specifying how the chunk should be indexed
 * @param f the read stream to pause while writing
 * @param processing an AtomicInteger keeping the number of chunks currently
 * being written (should be initialized to <code>0</code> the first time this
 * method is called)
 * @return a single that will emit exactly one item when the
 * operation has finished
 */
private Single<Void> addToStoreWithPause(Result<? extends ChunkMeta> chunk,
    String layer, IndexMeta indexMeta, ReadStream<Buffer> f, AtomicInteger processing) {
  // pause stream while chunk is being written
  f.pause();
  
  // count number of chunks being written
  processing.incrementAndGet();

  return addToStore(chunk.getChunk(), chunk.getMeta(), layer, indexMeta)
      .doOnSuccess(v -> {
        // resume stream only after all chunks from the current
        // buffer have been stored
        if (processing.decrementAndGet() == 0) {
          f.resume();
        }
      });
}
 
开发者ID:georocket,项目名称:georocket,代码行数:35,代码来源:ImporterVerticle.java


示例7: handleSingleResponseBody

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
default <R> void handleSingleResponseBody(RoutingContext ctx, Class<R> clazz,  R ret) {
    HttpServerResponse resp = ctx.response();
    if (ret != null) {
        if (ret instanceof String) {
            resp.end((String) ret);
        } else if (ret instanceof Buffer) {
            resp.end((Buffer) ret);
        } else if (ret instanceof JsonObject) {
            resp.end(((JsonObject) ret).encode());
        } else if (ret instanceof JsonArray) {
            resp.end(((JsonArray) ret).encode());
        } else {
            resp.end(objectToString(ret));
        }
    }
}
 
开发者ID:egopulse,项目名称:java-code-gen,代码行数:17,代码来源:RouteRegistrarHelper.java


示例8: delayToObservable

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
public void delayToObservable(HttpServer server) {
  server.requestHandler(request -> {
    if (request.method() == HttpMethod.POST) {

      // Stop receiving buffers
      request.pause();

      checkAuth(res -> {

        // Now we can receive buffers again
        request.resume();

        if (res.succeeded()) {
          Observable<Buffer> observable = request.toObservable();
          observable.subscribe(buff -> {
            // Get buffers
          });
        }
      });
    }
  });
}
 
开发者ID:vert-x3,项目名称:vertx-rx,代码行数:23,代码来源:RxifiedExamples.java


示例9: testHttpClient

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Test
public void testHttpClient() {
  HttpServer server = vertx.createHttpServer(new HttpServerOptions().setPort(8080));
  server.requestStream().handler(req -> {
    req.response().setChunked(true).end("some_content");
  });
  try {
    server.listen(ar -> {
      HttpClient client = vertx.createHttpClient(new HttpClientOptions());
      client.request(HttpMethod.GET, 8080, "localhost", "/the_uri", resp -> {
        Buffer content = Buffer.buffer();
        Observable<Buffer> observable = resp.toObservable();
        observable.forEach(content::appendBuffer, err -> fail(), () -> {
          assertEquals("some_content", content.toString("UTF-8"));
          testComplete();
        });
      }).end();
    });
    await();
  } finally {
    server.close();
  }
}
 
开发者ID:vert-x3,项目名称:vertx-rx,代码行数:24,代码来源:CoreApiTest.java


示例10: testWebsocketClient

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Test
public void testWebsocketClient() {
  HttpServer server = vertx.createHttpServer(new HttpServerOptions().setPort(8080));
  server.websocketStream().handler(ws -> {
    ws.write(Buffer.buffer("some_content"));
    ws.close();
  });
  server.listen(ar -> {
    HttpClient client = vertx.createHttpClient(new HttpClientOptions());
    client.websocket(8080, "localhost", "/the_uri", ws -> {
      Buffer content = Buffer.buffer();
      Observable<Buffer> observable = ws.toObservable();
      observable.forEach(content::appendBuffer, err -> fail(), () -> {
        server.close();
        assertEquals("some_content", content.toString("UTF-8"));
        testComplete();
      });
    });
  });
  await();
}
 
开发者ID:vert-x3,项目名称:vertx-rx,代码行数:22,代码来源:CoreApiTest.java


示例11: testWebsocketClientFlatMap

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Test
public void testWebsocketClientFlatMap() {
  HttpServer server = vertx.createHttpServer(new HttpServerOptions().setPort(8080));
  server.websocketStream().handler(ws -> {
    ws.write(Buffer.buffer("some_content"));
    ws.close();
  });
  server.listen(ar -> {
    HttpClient client = vertx.createHttpClient(new HttpClientOptions());
    Buffer content = Buffer.buffer();
    client.
        websocketStream(8080, "localhost", "/the_uri").
        toObservable().
        flatMap(WebSocket::toObservable).
        forEach(content::appendBuffer, err -> fail(), () -> {
          server.close();
          assertEquals("some_content", content.toString("UTF-8"));
          testComplete();
        });
  });
  await();
}
 
开发者ID:vert-x3,项目名称:vertx-rx,代码行数:23,代码来源:CoreApiTest.java


示例12: testGet

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Test
public void testGet() {
  int times = 5;
  waitFor(times);
  HttpServer server = vertx.createHttpServer(new HttpServerOptions().setPort(8080));
  server.requestStream().handler(req -> req.response().setChunked(true).end("some_content"));
  try {
    server.listen(ar -> {
      client = WebClient.wrap(vertx.createHttpClient(new HttpClientOptions()));
      Single<HttpResponse<Buffer>> single = client
        .get(8080, "localhost", "/the_uri")
        .as(BodyCodec.buffer())
        .rxSend();
      for (int i = 0; i < times; i++) {
        single.subscribe(resp -> {
          Buffer body = resp.body();
          assertEquals("some_content", body.toString("UTF-8"));
          complete();
        }, this::fail);
      }
    });
    await();
  } finally {
    server.close();
  }
}
 
开发者ID:vert-x3,项目名称:vertx-web,代码行数:27,代码来源:RxTest.java


示例13: testPost

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Test
public void testPost() {
  int times = 5;
  waitFor(times);
  HttpServer server = vertx.createHttpServer(new HttpServerOptions().setPort(8080));
  server.requestStream().handler(req -> req.bodyHandler(buff -> {
    assertEquals("onetwothree", buff.toString());
    req.response().end();
  }));
  try {
    server.listen(ar -> {
      client = WebClient.wrap(vertx.createHttpClient(new HttpClientOptions()));
      Observable<Buffer> stream = Observable.just(Buffer.buffer("one"), Buffer.buffer("two"), Buffer.buffer("three"));
      Single<HttpResponse<Buffer>> single = client
        .post(8080, "localhost", "/the_uri")
        .rxSendStream(stream);
      for (int i = 0; i < times; i++) {
        single.subscribe(resp -> complete(), this::fail);
      }
    });
    await();
  } finally {
    server.close();
  }
}
 
开发者ID:vert-x3,项目名称:vertx-web,代码行数:26,代码来源:RxTest.java


示例14: testResponseMissingBody

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Test
public void testResponseMissingBody() throws Exception {
  int times = 5;
  waitFor(times);
  HttpServer server = vertx.createHttpServer(new HttpServerOptions().setPort(8080));
  server.requestStream().handler(req -> req.response().setStatusCode(403).end());
  try {
    server.listen(ar -> {
      client = WebClient.wrap(vertx.createHttpClient(new HttpClientOptions()));
      Single<HttpResponse<Buffer>> single = client
        .get(8080, "localhost", "/the_uri")
        .rxSend();
      for (int i = 0; i < times; i++) {
        single.subscribe(resp -> {
          assertEquals(403, resp.statusCode());
          assertNull(resp.body());
          complete();
        }, this::fail);
      }
    });
    await();
  } finally {
    server.close();
  }
}
 
开发者ID:vert-x3,项目名称:vertx-web,代码行数:26,代码来源:RxTest.java


示例15: processFile

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
private static List<Tuple2<String, String>> processFile(final Vertx vertx, String inputFileStr) {
	File file = new File(inputFileStr);
	String fileName = inputFileStr.replaceFirst(".*/(\\w+).*","$1");
	List<Tuple2<String, String>> rules = new ArrayList<Tuple2<String, String>>();

	if (!file.isFile()) {
		final List<String> filesList = vertx.fileSystem().readDirBlocking(inputFileStr);

		for (final String dirFileStr : filesList) {
			List<Tuple2<String, String>> childRules = processFile(vertx, dirFileStr); // use directory name as rulegroup
			rules.addAll(childRules);
		}
		return rules;
	} else {
		Buffer buf = vertx.fileSystem().readFileBlocking(inputFileStr);
		try {
			final String ruleText = buf.toString();

			Tuple2<String, String> rule = (Tuple.of(fileName, ruleText));
			System.out.println("Loading in Rule:" + rule._1 + " of "+ inputFileStr);
			rules.add(rule);

			return rules;
		} catch (final DecodeException dE) {

		}

	}
	return null;
}
 
开发者ID:genny-project,项目名称:rulesservice,代码行数:31,代码来源:RulesLoader.java


示例16: get

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
private Single<String> get(Vertx vertx, URI uri){
	WebClient client = WebClient.create(vertx);
	Single<HttpResponse<Buffer>> responseHandler = 
			client.get(uri.getPort(), uri.getHost(), uri.getPath()).rxSend();

	return responseHandler.map(response -> response.body().toString());
}
 
开发者ID:FroMage,项目名称:redpipe,代码行数:8,代码来源:HelloResource.java


示例17: writeTo

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@Override
public void writeTo(Buffer buffer, Class<?> type, Type genericType, Annotation[] annotations, MediaType mediaType,
		MultivaluedMap<String, Object> httpHeaders, OutputStream entityStream)
		throws IOException, WebApplicationException {
	for(int i=0;i<buffer.length();i++)
		entityStream.write(buffer.getByte(i));
}
 
开发者ID:FroMage,项目名称:redpipe,代码行数:8,代码来源:BufferBodyWriter.java


示例18: observeUserExists

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
private Observable<Boolean> observeUserExists(String username) {
    String maybeJsonUser = usersCache.get(username);
    if (maybeJsonUser != null) {
        return Observable.just(true);
    }

    ObservableHandler<HttpClientResponse> responseObservable = RxHelper.observableHandler();
    client.getNow(userEndpoint(username), responseObservable.toHandler());
    return responseObservable
        .flatMap(response -> {
            if (response.statusCode() == HTTP_OK) {
                ObservableHandler<Buffer> observable = RxHelper.observableHandler();
                response.bodyHandler(observable.toHandler());
                return observable;
            } else {
                throw new RuntimeException(String.format("User call returned %d status",
                        response.statusCode()));
            }                
        })
        .map(buffer -> Json.decodeValue(buffer.toString(UTF_8.name()), ObjectNode.class))
        .map(user -> {
            boolean exists = user.has("user");
            if (exists) {
                // FIXME: not caching "non existing" users.
                usersCache.put(username, cacheSerialize(user));
            }
            return exists;
        });
}
 
开发者ID:MiguelGL,项目名称:popular-purchases-demo,代码行数:30,代码来源:RecentPurchasesHandler.java


示例19: importXML

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
/**
 * Imports an XML file from the given input stream into the store
 * @param f the XML file to read
 * @param correlationId a unique identifier for this import process
 * @param filename the name of the file currently being imported
 * @param timestamp denotes when the import process has started
 * @param layer the layer where the file should be stored (may be null)
 * @param tags the list of tags to attach to the file (may be null)
 * @param properties the map of properties to attach to the file (may be null)
 * @param fallbackCRSString the CRS which should be used if the imported
 * file does not specify one (may be <code>null</code>)
 * @return a single that will emit when the file has been imported
 */
protected Single<Integer> importXML(ReadStream<Buffer> f, String correlationId,
    String filename, long timestamp, String layer, List<String> tags,
    Map<String, Object> properties, String fallbackCRSString) {
  UTF8BomFilter bomFilter = new UTF8BomFilter();
  Window window = new Window();
  XMLSplitter splitter = new FirstLevelSplitter(window);
  AtomicInteger processing = new AtomicInteger(0);
  XMLCRSIndexer crsIndexer = new XMLCRSIndexer();
  return f.toObservable()
      .map(buf -> (io.vertx.core.buffer.Buffer)buf.getDelegate())
      .map(bomFilter::filter)
      .doOnNext(window::append)
      .lift(new XMLParserOperator())
      .doOnNext(e -> {
        // save the first CRS found in the file
        if (crsIndexer.getCRS() == null) {
          crsIndexer.onEvent(e);
        }
      })
      .flatMap(splitter::onEventObservable)
      .flatMapSingle(result -> {
        String crsString = fallbackCRSString;
        if (crsIndexer.getCRS() != null) {
          crsString = crsIndexer.getCRS();
        }
        IndexMeta indexMeta = new IndexMeta(correlationId, filename,
            timestamp, tags, properties, crsString);
        return addToStoreWithPause(result, layer, indexMeta, f, processing);
      })
      .count()
      .toSingle();
}
 
开发者ID:georocket,项目名称:georocket,代码行数:46,代码来源:ImporterVerticle.java


示例20: getReqBody

import io.vertx.rxjava.core.buffer.Buffer; //导入依赖的package包/类
@SuppressWarnings("unchecked")
default <T> T getReqBody(Class<T> clazz, RoutingContext ctx) {
    if (String.class.equals(clazz)) {
        return (T) ctx.getBodyAsString();
    }
    if (JsonObject.class.equals(clazz)) {
        return (T) ctx.getBodyAsJson();
    }
    if (Buffer.class.equals(clazz)) {
        return (T) ctx.getBody();
    }
    return stringToObject(clazz, ctx.getBodyAsString());
}
 
开发者ID:egopulse,项目名称:java-code-gen,代码行数:14,代码来源:RouteRegistrarHelper.java



注:本文中的io.vertx.rxjava.core.buffer.Buffer类示例整理自Github/MSDocs等源码及文档管理平台,相关代码片段筛选自各路编程大神贡献的开源项目,源码版权归原作者所有,传播和使用请参考对应项目的License;未经允许,请勿转载。


鲜花

握手

雷人

路过

鸡蛋
该文章已有0人参与评论

请发表评论

全部评论

专题导读
上一篇:
Java InvalidApplicationException类代码示例发布时间:2022-05-22
下一篇:
Java PeerDiscovery类代码示例发布时间:2022-05-22
热门推荐
阅读排行榜

扫描微信二维码

查看手机版网站

随时了解更新最新资讯

139-2527-9053

在线客服(服务时间 9:00~18:00)

在线QQ客服
地址:深圳市南山区西丽大学城创智工业园
电邮:jeky_zhao#qq.com
移动电话:139-2527-9053

Powered by 互联科技 X3.4© 2001-2213 极客世界.|Sitemap