- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.errors.SerializationException;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
class MultithreadingDemo extends Thread
{
public void run()
{
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "xxx:443");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
io.confluent.kafka.serializers.KafkaAvroSerializer.class);
props.put("schema.registry.url", "xxx");
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put("security.protocol", "SSL");
props.put("ssl.truststore.location", "xxx");
props.put("ssl.truststore.password", "xxx");
props.put("ssl.keystore.location", "xxx");
props.put("ssl.keystore.password", "xxx");
props.put("ssl.key.password", "xxx");
KafkaProducer producer = new KafkaProducer(props);
String userSchema = "{ \"name\": \"MyClass\", \"type\": \"record\", \"namespace\":
\"com.oop.hts\", \"fields\": [ { \"name\": \"appId\", \"type\":
\"string\" }, { \"name\": \"appName\", \"type\": \"string\" }, {
\"name\": \"groups\", \"type\": \"string\" }, { \"name\": \"subGroups\",
\"type\": \"string\" }, { \"name\": \"jobType\", \"type\": \"string\"
}, { \"name\": \"appStartTime\", \"type\": \"string\" }, {
\"name\": \"appEndTime\", \"type\": \"string\" }, { \"name\":
\"appDuration\", \"type\": \"int\" }, { \"name\": \"cpuTime\",
\"type\": \"int\" }, { \"name\": \"runTime\", \"type\": \"int\" },
{ \"name\": \"memoryUsage\", \"type\": \"int\" }, { \"name\":
\"appStatus\", \"type\": \"string\" }, { \"name\": \"appResult\",
\"type\": \"string\" }, { \"name\": \"failureREason\", \"type\":
\"string\" }, { \"name\": \"recordCount\", \"type\": \"string\" },
{ \"name\": \"numexecutors\", \"type\": \"string\" }, { \"name\":
\"executorcores\", \"type\": \"string\" }, { \"name\":
\"executormemory\", \"type\": \"string\" } ] }\n" +
"";
System.out.println("schema:" + userSchema);
Schema.Parser parser = new Schema.Parser();
Schema schema = parser.parse(userSchema);
GenericRecord avroRecord = new GenericData.Record(schema);
//avroRecord.put("f1", "value777");
System.out.println("----" + avroRecord);
avroRecord.put("appId","spark-d0731a81f1b64f109c5d985c1b2e0011");
avroRecord.put("appName","H@S-UCR");
avroRecord.put("groups","");
avroRecord.put("subGroups","");
avroRecord.put("jobType","");
avroRecord.put("appStartTime","2020-04-13T10:02:25.902");
avroRecord.put("appEndTime","2020-04-13T10:02:25.902");
avroRecord.put("appDuration",4110);
avroRecord.put("cpuTime",337468);
avroRecord.put("runTime",1198987);
avroRecord.put("memoryUsage",234933352);
avroRecord.put("appStatus","Running");
avroRecord.put("appResult","InProgress");
avroRecord.put("failureREason","");
avroRecord.put("recordCount","0");
avroRecord.put("numexecutors","25");
avroRecord.put("executorcores","15");
avroRecord.put("executormemory","60g");
System.out.println("----"+ avroRecord);
ProducerRecord<String, GenericRecord> record = new ProducerRecord<String,
GenericRecord>("kaas.topic", avroRecord);
try {
producer.send(record);
System.out.println("Successfully produced the records to the Kafka topic :
kaas.dqhats.target ");
} catch(SerializationException e) {
System.out.println("An Exception occured" + e.getMessage());
e.printStackTrace();
}
}
}
// Main Class
public class Multithread
{
public static void main(String[] args)
{
int n = 8; // Number of threads
for (int i=0; i<n; i++)
{
MultithreadingDemo object = new MultithreadingDemo();
object.start();
}
}
}
我想使用多线程向 kafka 分区生成多条消息。(这是检查 kafka 主题/分区性能/容量所需的)
使用以下代码,我无法并行向 kafka 分区生成消息。
寻求帮助。
使用多线程同时向 Kafka 分区发布多条消息以进行测试以检查性能
任何人都可以帮助我使用多线程同时向 Kafka 分区发布多条消息。
最佳答案
send()
方法只会将消息放入缓冲区中,并且消息将作为单独线程的一部分发送。本质上,这就是所展示的生产者的异步本质。
此外,在调用 send()
方法之后,此调用返回的 Future
对象将被忽略,因此您实际上无法知道是否您的消息是否已发送。
您可以尝试:
生产者.send(record).get();
这将在继续之前等待 Kafka 的响应,如果将该消息发送到 Kafka 时出现任何问题,您将收到错误消息。
或者
send()
之后调用 flush()
方法。顾名思义,此方法将刷新缓冲区中的消息,但 here is the reference documentation为此,如果您想了解更多信息。
希望这有帮助!
关于java - 使用多线程同时向 Kafka 分区发布多条消息以进行测试以检查性能,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/61682493/
我获得了一些源代码示例,我想测试一些功能。不幸的是,我在执行程序时遇到问题: 11:41:31 [linqus@ottsrvafq1 example]$ javac -g test/test.jav
我想测试ggplot生成的两个图是否相同。一种选择是在绘图对象上使用all.equal,但我宁愿进行更艰巨的测试以确保它们相同,这似乎是identical()为我提供的东西。 但是,当我测试使用相同d
我确实使用 JUnit5 执行我的 Maven 测试,其中所有测试类都有 @ExtendWith({ProcessExtension.class}) 注释。如果是这种情况,此扩展必须根据特殊逻辑使测试
在开始使用 Node.js 开发有用的东西之前,您的流程是什么?您是否在 VowJS、Expresso 上创建测试?你使用 Selenium 测试吗?什么时候? 我有兴趣获得一个很好的工作流程来开发我
这个问题已经有答案了: What is a NullPointerException, and how do I fix it? (12 个回答) 已关闭 3 年前。 基于示例here ,我尝试为我的
我正在考虑测试一些 Vue.js 组件,作为 Laravel 应用程序的一部分。所以,我有一个在 Blade 模板中使用并生成 GET 的组件。在 mounted 期间请求生命周期钩子(Hook)。假
考虑以下程序: #include struct Test { int a; }; int main() { Test t=Test(); std::cout<
我目前的立场是:如果我使用 web 测试(在我的例子中可能是通过 VS.NET'08 测试工具和 WatiN)以及代码覆盖率和广泛的数据来彻底测试我的 ASP.NET 应用程序,我应该不需要编写单独的
我正在使用 C#、.NET 4.7 我有 3 个字符串,即。 [test.1, test.10, test.2] 我需要对它们进行排序以获得: test.1 test.2 test.10 我可能会得到
我有一个 ID 为“rv_list”的 RecyclerView。单击任何 RecyclerView 项目时,每个项目内都有一个可见的 id 为“star”的 View 。 我想用 expresso
我正在使用 Jest 和模拟器测试 Firebase 函数,尽管这些测试可能来自竞争条件。所谓 flakey,我的意思是有时它们会通过,有时不会,即使在同一台机器上也是如此。 测试和函数是用 Type
我在测试我与 typeahead.js ( https://github.com/angular-ui/bootstrap/blob/master/src/typeahead/typeahead.js
我正在尝试使用 Teamcity 自动运行测试,但似乎当代理编译项目时,它没有正确完成,因为当我运行运行测试之类的命令时,我收到以下错误: fatal error: 'Pushwoosh/PushNo
这是我第一次玩 cucumber ,还创建了一个测试和 API 的套件。我的问题是在测试 API 时是否需要运行它? 例如我脑子里有这个, 启动 express 服务器作为后台任务 然后当它启动时(我
我有我的主要应用程序项目,然后是我的测试的第二个项目。将所有类型的测试存储在该测试项目中是一种好的做法,还是应该将一些测试驻留在主应用程序项目中? 我应该在我的主项目中保留 POJO JUnit(测试
我正在努力弄清楚如何实现这个计数。模型是用户、测试、等级 用户 has_many 测试,测试 has_many 成绩。 每个等级都有一个计算分数(strong_pass、pass、fail、stron
我正在尝试测试一些涉及 OkHttp3 的下载代码,但不幸失败了。目标:测试 下载图像文件并验证其是否有效。平台:安卓。此代码可在生产环境中运行,但测试代码没有任何意义。 产品代码 class Fil
当我想为 iOS 运行 UI 测试时,我收到以下消息: SetUp : System.Exception : Unable to determine simulator version for X 堆
我正在使用 Firebase Remote Config 在 iOS 上设置 A/B 测试。 一切都已设置完毕,我正在 iOS 应用程序中读取服务器端默认值。 但是在多个模拟器上尝试,它们都读取了默认
[已编辑]:我已经用 promise 方式更改了我的代码。 我正在写 React with this starter 由 facebook 创建,我是测试方面的新手。 现在我有一个关于图像的组件,它有
我是一名优秀的程序员,十分优秀!