(systemTestEnv, testcaseEnv, onlyThisEntityId=None)
| 1688 | |
| 1689 | |
| 1690 | def start_migration_tool(systemTestEnv, testcaseEnv, onlyThisEntityId=None): |
| 1691 | clusterConfigList = systemTestEnv.clusterEntityConfigDictList |
| 1692 | migrationToolConfigList = system_test_utils.get_dict_from_list_of_dicts(clusterConfigList, "role", "migration_tool") |
| 1693 | |
| 1694 | for migrationToolConfig in migrationToolConfigList: |
| 1695 | |
| 1696 | entityId = migrationToolConfig["entity_id"] |
| 1697 | |
| 1698 | if onlyThisEntityId is None or entityId == onlyThisEntityId: |
| 1699 | |
| 1700 | host = migrationToolConfig["hostname"] |
| 1701 | jmxPort = migrationToolConfig["jmx_port"] |
| 1702 | role = migrationToolConfig["role"] |
| 1703 | kafkaHome = system_test_utils.get_data_by_lookup_keyval(clusterConfigList, "entity_id", entityId, "kafka_home") |
| 1704 | javaHome = system_test_utils.get_data_by_lookup_keyval(clusterConfigList, "entity_id", entityId, "java_home") |
| 1705 | jmxPort = system_test_utils.get_data_by_lookup_keyval(clusterConfigList, "entity_id", entityId, "jmx_port") |
| 1706 | kafkaRunClassBin = kafkaHome + "/bin/kafka-run-class.sh" |
| 1707 | |
| 1708 | logger.info("starting kafka migration tool", extra=d) |
| 1709 | migrationToolLogPath = get_testcase_config_log_dir_pathname(testcaseEnv, "migration_tool", entityId, "default") |
| 1710 | migrationToolLogPathName = migrationToolLogPath + "/migration_tool.log" |
| 1711 | testcaseEnv.userDefinedEnvVarDict["migrationToolLogPathName"] = migrationToolLogPathName |
| 1712 | |
| 1713 | testcaseConfigsList = testcaseEnv.testcaseConfigsList |
| 1714 | numProducers = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "num.producers") |
| 1715 | numStreams = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "num.streams") |
| 1716 | producerConfig = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "producer.config") |
| 1717 | consumerConfig = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "consumer.config") |
| 1718 | zkClientJar = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "zkclient.01.jar") |
| 1719 | kafka07Jar = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "kafka.07.jar") |
| 1720 | whiteList = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "whitelist") |
| 1721 | logFile = system_test_utils.get_data_by_lookup_keyval(testcaseConfigsList, "entity_id", entityId, "log_filename") |
| 1722 | |
| 1723 | cmdList = ["ssh " + host, |
| 1724 | "'JAVA_HOME=" + javaHome, |
| 1725 | "JMX_PORT=" + jmxPort, |
| 1726 | kafkaRunClassBin + " kafka.tools.KafkaMigrationTool", |
| 1727 | "--whitelist=" + whiteList, |
| 1728 | "--num.producers=" + numProducers, |
| 1729 | "--num.streams=" + numStreams, |
| 1730 | "--producer.config=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + producerConfig, |
| 1731 | "--consumer.config=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + consumerConfig, |
| 1732 | "--zkclient.01.jar=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + zkClientJar, |
| 1733 | "--kafka.07.jar=" + systemTestEnv.SYSTEM_TEST_BASE_DIR + "/" + kafka07Jar, |
| 1734 | " &> " + migrationToolLogPath + "/migrationTool.log", |
| 1735 | " & echo pid:$! > " + migrationToolLogPath + "/entity_" + entityId + "_pid'"] |
| 1736 | |
| 1737 | cmdStr = " ".join(cmdList) |
| 1738 | logger.debug("executing command: [" + cmdStr + "]", extra=d) |
| 1739 | system_test_utils.async_sys_call(cmdStr) |
| 1740 | time.sleep(5) |
| 1741 | |
| 1742 | pidCmdStr = "ssh " + host + " 'cat " + migrationToolLogPath + "/entity_" + entityId + "_pid' 2> /dev/null" |
| 1743 | logger.debug("executing command: [" + pidCmdStr + "]", extra=d) |
| 1744 | subproc = system_test_utils.sys_call_return_subproc(pidCmdStr) |
| 1745 | |
| 1746 | # keep track of the remote entity pid in a dictionary |
| 1747 | for line in subproc.stdout.readlines(): |
nothing calls this directly
no test coverage detected