Подпрограммы и каналы − учебник
В этом руководстве вы узнаете, как использовать подпрограммы в IntelliJ IDEA для выполнения сетевых запросов без блокировки основной нити или обратных вызовов.
Вы узнаете:
Зачем и как использовать приостанавливаемые функции для выполнения сетевых запросов.
Как отправлять запросы параллельно с помощью подпрограмм.
Как обмениваться данными между различными подпрограммами с помощью каналов.
Для сетевых запросов вам потребуется библиотека Retrofit, но подход, показанный в этом руководстве, аналогичным образом работает с другими библиотеками, поддерживающими подпрограммы.
Перед началом
Загрузите и установите последнюю версию IntelliJ IDEA.
-
Скопируйте шаблон проекта шаблон проекта, выбрав Получить из VCS на экране «Добро пожаловать» или выбрав Файл | Новый | Проект из системы управления версиями.
Вы также можете скопировать его из командной строки:
git clone https://github.com/kotlin-hands-on/intro-coroutines
Создайте токен разработчика GitHub
В проекте вы будете использовать API GitHub. Для получения доступа укажите имя вашей учетной записи GitHub и пароль или токен. Если у вас включена двухфакторная аутентификация, токена будет достаточно.
Создайте новый токен GitHub для использования API GitHub с вашей учетной записью:
-
Укажите имя вашего токена, например,
coroutines-tutorial:
Не выбирайте никакие области. Нажмите Создать токен в нижней части страницы.
Скопируйте сгенерированный токен.
Запуск кода
Программа загружает авторов для всех репозиториев в данной организации (по умолчанию «kotlin»). Позже вы добавите логику для сортировки пользователей по количеству их вкладов.
-
Откройте файл
src/contributors/main.ktи запустите функциюmain(). Вы увидите следующее окно:Если шрифт слишком мелкий, измените его, изменив значение
setDefaultFontSize(18f)в функцииmain(). Укажите имя пользователя GitHub и токен (или пароль) в соответствующих полях.
Убедитесь, что в раскрывающемся списке «Вариант» выбран параметр БЛОКИРУЮЩИЙ.
Нажмите Загрузить авторов. Интерфейс пользователя должен на некоторое время заблокироваться, а затем отобразить список авторов.
Откройте вывод программы, чтобы убедиться, что данные загружены. Список авторов регистрируется после каждого успешного запроса.
Существуют разные способы реализации этой логики: с помощью блокирующих запросов или обратных вызовов. Вы сравните эти решения с одним, использующим подпрограммы, и увидите, как каналы можно использовать для обмена данными между различными подпрограммами.
Блокирующие запросы
Для выполнения HTTP-запросов к GitHub вы будете использовать библиотеку Retrofit. Она позволяет запрашивать список репозиториев в рамках данной организации и список авторов каждого репозитория:
interface GitHubService {
@GET("orgs/{org}/repos?per_page=100")
fun getOrgReposCall(
@Path("org") org: String
): Call<List<Repo>>
@GET("repos/{owner}/{repo}/contributors?per_page=100")
fun getRepoContributorsCall(
@Path("owner") owner: String,
@Path("repo") repo: String
): Call<List<User>>
}
Этот API используется функцией loadContributorsBlocking() для получения списка авторов для данной организации.
-
Откройте
src/tasks/Request1Blocking.ktдля просмотра её реализации:fun loadContributorsBlocking(service: GitHubService, req: RequestData): List<User> { val repos = service .getOrgReposCall(req.org) // #1 .execute() // #2 .also { logRepos(req, it) } // #3 .body() ?: emptyList() // #4 return repos.flatMap { repo -> service .getRepoContributorsCall(req.org, repo.name) // #1 .execute() // #2 .also { logUsers(repo, it) } // #3 .bodyList() // #4 }.aggregate() }Сначала вы получаете список репозиториев для данной организации и сохраняете его в списке
repos. Затем для каждого репозитория запрашивается список авторов, и все списки объединяются в один итоговый список авторов.getOrgReposCall()иgetRepoContributorsCall()оба возвращают экземпляр класса*Call(#1). На этом этапе ни один запрос не отправляется.Затем вызывается
*Call.execute()для выполнения запроса (#2).execute()— это синхронный вызов, который блокирует подлежащую нить.При получении ответа результат регистрируется с помощью вызова конкретных функций
logRepos()иlogUsers()(#3). Если HTTP-ответ содержит ошибку, эта ошибка будет зарегистрирована здесь.Наконец, получите тело ответа, которое содержит необходимые данные. Для этого учебника вы будете использовать пустой список в качестве результата в случае ошибки и будете регистрировать соответствующую ошибку (
#4).
-
Чтобы избежать повторения
.body() ?: emptyList(), объявляется функция-расширениеbodyList().fun <T> Response<List<T>>.bodyList(): List<T> { return body() ?: emptyList() } -
Запустите программу ещё раз и посмотрите вывод в IntelliJ IDEA. Он должен выглядеть примерно так:
1770 [AWT-EventQueue-0] INFO Contributors - kotlin: loaded 40 repos 2025 [AWT-EventQueue-0] INFO Contributors - kotlin-examples: loaded 23 contributors 2229 [AWT-EventQueue-0] INFO Contributors - kotlin-koans: loaded 45 contributors ...
Первым элементом на каждой строке является количество миллисекунд, прошедших с момента запуска программы, затем имя потока в квадратных скобках. Вы можете видеть, из какого потока вызывается запрос загрузки.
Последним элементом на каждой строке является фактическое сообщение: сколько репозиториев или авторов было загружено.
Этот лог-вывод демонстрирует, что все результаты были зарегистрированы из главного потока. При запуске кода с опцией БЛОКИРУЮЩИЙ окно замирает и не реагирует на ввод, пока загрузка не завершится. Все запросы выполняются из того же потока, что и функция
loadContributorsBlocking(), которая является главным потоком пользовательского интерфейса (в Swing это поток диспетчеризации событий AWT). Этот главный поток блокируется, и поэтому пользовательский интерфейс замирает:После загрузки списка авторов результат обновляется.
-
В
src/contributors/Contributors.kt, найдите функциюloadContributors(), отвечающую за выбор способа загрузки авторов, и посмотрите, как вызываетсяloadContributorsBlocking().when (getSelectedVariant()) { BLOCKING -> { // Blocking UI thread val users = loadContributorsBlocking(service, req) updateResults(users, startTime) } }Вызов
updateResults()происходит сразу после вызоваloadContributorsBlocking().updateResults()обновляет пользовательский интерфейс, поэтому его всегда следует вызывать из потока пользовательского интерфейса.Поскольку
loadContributorsBlocking()также вызывается из потока пользовательского интерфейса, поток пользовательского интерфейса блокируется, и пользовательский интерфейс замирает.
Задача 1
Первая задача поможет вам ознакомиться с задачей. В настоящее время имя каждого автора повторяется несколько раз, по одному разу для каждого проекта, в котором он принимал участие. Реализуйте функцию aggregate(), объединяющую пользователей таким образом, чтобы каждый автор добавлялся только один раз. Свойство User.contributions должно содержать общее количество вкладов данного пользователя во всех проектах. Результирующий список должен быть отсортирован по убыванию количества вкладов.
Откройте src/tasks/Aggregation.kt и реализуйте функцию List<User>.aggregate(). Пользователи должны быть отсортированы по общему количеству их вкладов.
Соответствующий тестовый файл test/tasks/AggregationKtTest.kt демонстрирует пример ожидаемого результата.
После выполнения этой задачи результирующий список для организации "kotlin" должен быть похож на следующий:

Решение задачи 1
Для группировки пользователей по логину используйте
groupBy(), которая возвращает карту, связывающую логин со всеми экземплярами пользователя с этим логином в разных репозиториях.Для каждой записи в карте подсчитайте общее количество вкладов для каждого пользователя и создайте новый экземпляр класса
Userс заданным именем и общим количеством вкладов.-
Отсортируйте результирующий список по убыванию:
fun List<User>.aggregate(): List<User> = groupBy { it.login } .map { (login, group) -> User(login, group.sumOf { it.contributions }) } .sortedByDescending { it.contributions }
Альтернативным решением является использование функции groupingBy() вместо groupBy().
Обработчики событий
Предыдущее решение работает, но оно блокирует поток и, следовательно, замораживает пользовательский интерфейс. Традиционный подход, который избегает этого, заключается в использовании обработчиков событий.
Вместо вызова кода, который должен быть вызван сразу после завершения операции, вы можете вынести его в отдельный обработчик событий, часто лямбду, и передать эту лямбду вызывающему объекту, чтобы он был вызван позже.
Чтобы сделать пользовательский интерфейс отзывчивым, вы можете либо перенести все вычисления в отдельный поток, либо переключиться на API Retrofit, который использует обработчики событий вместо блокирующих вызовов.
Использование фонового потока
-
Откройте
src/tasks/Request2Background.ktи ознакомьтесь с его реализацией. Сначала все вычисления перемещаются в другой поток. Функцияthread()запускает новый поток:thread { loadContributorsBlocking(service, req) }Теперь, когда вся загрузка перенесена в отдельный поток, основной поток свободен и может быть занят другими задачами:
-
Подпись функции
loadContributorsBackground()изменяется. Она принимает обработчикupdateResults()событий в качестве последнего аргумента, чтобы вызвать его после завершения всей загрузки:fun loadContributorsBackground( service: GitHubService, req: RequestData, updateResults: (List<User>) -> Unit ) -
Теперь, когда вызывается
loadContributorsBackground(), вызовupdateResults()передается в обработчик событий, а не сразу после него, как это было раньше:loadContributorsBackground(service, req) { users -> SwingUtilities.invokeLater { updateResults(users, startTime) } }Вызывая
SwingUtilities.invokeLater, вы гарантируете, что вызовupdateResults(), обновляющий результаты, происходит в главном потоке пользовательского интерфейса (поток диспетчеризации событий AWT).
Однако, если вы попытаетесь загрузить участников через опцию BACKGROUND, вы увидите, что список обновляется, но ничего не меняется.
Задача 2
Исправьте функцию loadContributorsBackground() в src/tasks/Request2Background.kt таким образом, чтобы результирующий список отображался в пользовательском интерфейсе.
Решение задачи 2
Если вы попытаетесь загрузить участников, вы увидите в журнале, что участники загрузились, но результат не отображается. Чтобы исправить это, вызовите updateResults() на результирующем списке пользователей:
thread {
updateResults(loadContributorsBlocking(service, req))
}
Убедитесь, что вы явно вызываете логику, переданную в обработчике событий. В противном случае ничего не произойдет.
Использование API обратного вызова Retrofit
В предыдущем решении вся логика загрузки перенесена в фоновый поток, но это всё ещё не наилучшее использование ресурсов. Все запросы загрузки выполняются последовательно, и поток блокируется, ожидая результата загрузки, в то время как он мог бы быть занят другими задачами. В частности, поток мог бы начать загрузку другого запроса, чтобы получить весь результат раньше.
Обработка данных для каждого репозитория должна быть разделена на две части: загрузка и обработка результирующего ответа. Вторая часть обработки должна быть вынесена в обработчик событий.
Загрузка для каждого репозитория может быть начата до получения результата для предыдущего репозитория (и вызова соответствующего обработчика событий):

API обратного вызова Retrofit может помочь в этом. Функция Call.enqueue() запускает HTTP-запрос и принимает обработчик событий в качестве аргумента. В этом обработчике вы должны указать, что нужно сделать после каждого запроса.
Откройте src/tasks/Request3Callbacks.kt и ознакомьтесь с реализацией loadContributorsCallbacks(), которая использует этот API:
fun loadContributorsCallbacks(
service: GitHubService, req: RequestData,
updateResults: (List<User>) -> Unit
) {
service.getOrgReposCall(req.org).onResponse { responseRepos -> // #1
logRepos(req, responseRepos)
val repos = responseRepos.bodyList()
val allUsers = mutableListOf<User>()
for (repo in repos) {
service.getRepoContributorsCall(req.org, repo.name)
.onResponse { responseUsers -> // #2
logUsers(repo, responseUsers)
val users = responseUsers.bodyList()
allUsers += users
}
}
}
// TODO: Why doesn't this code work? How to fix that?
updateResults(allUsers.aggregate())
}
Для удобства этот фрагмент кода использует функцию расширения
onResponse(), объявленную в том же файле. Она принимает лямбду в качестве аргумента вместо выражения объекта.Логика обработки ответов вынесена в обработчики событий: соответствующие лямбды начинаются на строках
#1и#2.
Однако, предложенное решение не работает. Если вы запустите программу и загрузите участников, выбрав опцию ОБРАБОТЧИКИ СОБЫТИЙ, то ничего не отобразится. Однако тесты, которые сразу возвращают результат, проходят.
Подумайте, почему предоставленный код не работает как ожидается, и попробуйте исправить его, или см. решения ниже.
Задача 3 (дополнительная)
Перепишите код в файле src/tasks/Request3Callbacks.kt так, чтобы загруженный список участников отображался.
Первое попытка решения задачи 3
В текущем решении многие запросы запускаются одновременно, что уменьшает общее время загрузки. Однако результат не загружается. Это происходит потому, что обработчик событий updateResults() вызывается сразу после запуска всех запросов на загрузку, до того, как список allUsers заполнится данными.
Вы могли бы попробовать исправить это с помощью такого изменения:
val allUsers = mutableListOf<User>()
for ((index, repo) in repos.withIndex()) { // #1
service.getRepoContributorsCall(req.org, repo.name)
.onResponse { responseUsers ->
logUsers(repo, responseUsers)
val users = responseUsers.bodyList()
allUsers += users
if (index == repos.lastIndex) { // #2
updateResults(allUsers.aggregate())
}
}
}
Сначала вы проходите по списку репозиториев с индексом (
#1).Затем из каждого обработчика событий вы проверяете, является ли это последняя итерация (
#2).И если это так, то результат обновляется.
Однако и этот код не решает поставленную задачу. Попробуйте найти решение самостоятельно, или см. решение ниже.
Второе попытка решения задачи 3
Поскольку запросы на загрузку запускаются одновременно, нет гарантии, что результат последнего придёт последним. Результаты могут приходить в любом порядке.
Таким образом, если вы сравниваете текущий индекс с lastIndex как условие завершения, вы рискуете потерять результаты для некоторых репозиториев.
Если запрос, обрабатывающий последний репозиторий, возвращается быстрее, чем некоторые предыдущие запросы (что, вероятно, произойдёт), все результаты запросов, которые занимают больше времени, будут потеряны.
Один из способов исправить это — ввести индекс и проверить, были ли уже обработаны все репозитории:
val allUsers = Collections.synchronizedList(mutableListOf<User>())
val numberOfProcessed = AtomicInteger()
for (repo in repos) {
service.getRepoContributorsCall(req.org, repo.name)
.onResponse { responseUsers ->
logUsers(repo, responseUsers)
val users = responseUsers.bodyList()
allUsers += users
if (numberOfProcessed.incrementAndGet() == repos.size) {
updateResults(allUsers.aggregate())
}
}
}
Этот код использует синхронизированный вариант списка и AtomicInteger(), потому что, вообще говоря, нет гарантии, что разные обработчики событий, обрабатывающие getRepoContributors() запросы, всегда будут вызываться из одного и того же потока.
Третье попытка решения задачи 3
Ещё лучшее решение — использовать класс CountDownLatch. Он хранит счётчик, инициализированный количеством репозиториев. Этот счётчик декрементируется после обработки каждого репозитория. Затем он ожидает, пока счётчик защелки не обнулится, прежде чем обновлять результаты:
val countDownLatch = CountDownLatch(repos.size)
for (repo in repos) {
service.getRepoContributorsCall(req.org, repo.name)
.onResponse { responseUsers ->
// processing repository
countDownLatch.countDown()
}
}
countDownLatch.await()
updateResults(allUsers.aggregate())
Результат затем обновляется из главного потока. Это более прямо, чем делегирование логики дочерним потокам.
После рассмотрения этих трёх попыток решения вы можете увидеть, что написание корректного кода с обработчиками событий нетривиально и подвержено ошибкам, особенно когда участвуют несколько подчинённых потоков и синхронизация.
Функции приостановки
Вы можете реализовать ту же логику, используя функции приостановки. Вместо возвращения Call<List<Repo>>, определите вызов API как функцию приостановки следующим образом:
interface GitHubService {
@GET("orgs/{org}/repos?per_page=100")
suspend fun getOrgRepos(
@Path("org") org: String
): List<Repo>
}
getOrgRepos()определена как функцияsuspend. Когда вы используете функцию приостановки для выполнения запроса, подлежащая нить не блокируется. Более подробные сведения о том, как это работает, будут представлены в последующих разделах.getOrgRepos()возвращает результат напрямую, вместо возвращенияCall. Если результат неудачный, выбрасывается исключение.
В качестве альтернативы Retrofit позволяет возвращать результат, обернутый в Response. В этом случае тело результата предоставляется, и можно вручную проверить ошибки. В этом руководстве используются версии, возвращающие Response.
В src/contributors/GitHubService.kt, добавьте следующие объявления в интерфейс GitHubService:
interface GitHubService {
// getOrgReposCall & getRepoContributorsCall declarations
@GET("orgs/{org}/repos?per_page=100")
suspend fun getOrgRepos(
@Path("org") org: String
): Response<List<Repo>>
@GET("repos/{owner}/{repo}/contributors?per_page=100")
suspend fun getRepoContributors(
@Path("owner") owner: String,
@Path("repo") repo: String
): Response<List<User>>
}
Задача 4
Ваша задача — изменить код функции, загружающей участников, чтобы использовать две новые функции приостановки, getOrgRepos() и getRepoContributors(). Новая функция loadContributorsSuspend() помечена как suspend для использования нового API.
Скопируйте реализацию
loadContributorsBlocking(), которая определена вsrc/tasks/Request1Blocking.kt, вloadContributorsSuspend(), которая определена вsrc/tasks/Request4Suspend.kt.Измените код так, чтобы использовались новые функции приостановки вместо функций, возвращающих
Call.Запустите программу, выбрав опцию ПРИОСТАНОВИТЬ, и убедитесь, что пользовательский интерфейс по-прежнему отзывчив во время выполнения запросов GitHub.
Решение задачи 4
Замените .getOrgReposCall(req.org).execute() на .getOrgRepos(req.org) и повторите ту же замену для второго запроса "участников":
suspend fun loadContributorsSuspend(service: GitHubService, req: RequestData): List<User> {
val repos = service
.getOrgRepos(req.org)
.also { logRepos(req, it) }
.bodyList()
return repos.flatMap { repo ->
service.getRepoContributors(req.org, repo.name)
.also { logUsers(repo, it) }
.bodyList()
}.aggregate()
}
loadContributorsSuspend()должна быть определена как функцияsuspend.Вам больше не нужно вызывать
execute, которая ранее возвращалаResponse, потому что теперь функции API возвращаютResponseнапрямую. Обратите внимание, что эта деталь специфична для библиотеки Retrofit. С другими библиотеками API будет отличаться, но концепция остаётся той же.
Сопрограммы
Код с функциями приостановки выглядит похожим на "блокирующий" вариант. Основное отличие от блокирующего варианта заключается в том, что вместо блокирования потока сопрограмма приостанавливается:
block -> suspend thread -> coroutine
Запуск новой сопрограммы
Если вы посмотрите, как loadContributorsSuspend() используется в src/contributors/Contributors.kt, вы увидите, что она вызывается внутри launch . launch — это функция библиотеки, которая принимает лямбда-выражение в качестве аргумента:
launch {
val users = loadContributorsSuspend(req)
updateResults(users, startTime)
}
Здесь launch запускает новую вычислительную задачу, ответственная за загрузку данных и отображение результатов. Вычислительная задача является приостанавливаемой — при выполнении сетевых запросов она приостанавливается и освобождает подлежащую нить. Когда сетевой запрос возвращает результат, вычислительная задача возобновляется.
Такая приостанавливаемая вычислительная задача называется сопрограммой. Таким образом, в данном случае launch запускает новую сопрограмму, отвечающую за загрузку данных и отображение результатов.
Сопрограммы работают поверх потоков и могут быть приостановлены. Когда сопрограмма приостановлена, соответствующая вычислительная задача приостанавливается, удаляется из потока и сохраняется в памяти. Между тем поток свободен для выполнения других задач:

Когда вычислительная задача готова к продолжению, она возвращается в поток (не обязательно в тот же).
В примере loadContributorsSuspend() каждый запрос "участников" теперь ожидает результат, используя механизм приостановки. Сначала отправляется новый запрос. Затем, ожидая ответа, вся сопрограмма "загрузить участников", запущенная функцией launch, приостанавливается.
Сопрограмма возобновляется только после получения соответствующего ответа:

Пока ответ ожидает получения, поток свободен для выполнения других задач. Пользовательский интерфейс остаётся отзывчивым, несмотря на все запросы, выполняемые в основном потоке пользовательского интерфейса:
-
Запустите программу, используя опцию ПРИОСТАНОВИТЬ. Лог подтверждает, что все запросы отправляются в основной поток пользовательского интерфейса:
2538 [AWT-EventQueue-0 @coroutine#1] INFO Contributors - kotlin: loaded 30 repos 2729 [AWT-EventQueue-0 @coroutine#1] INFO Contributors - ts2kt: loaded 11 contributors 3029 [AWT-EventQueue-0 @coroutine#1] INFO Contributors - kotlin-koans: loaded 45 contributors ... 11252 [AWT-EventQueue-0 @coroutine#1] INFO Contributors - kotlin-coroutines-workshop: loaded 1 contributors
-
В логе вы можете увидеть, в какой сопрограмме выполняется соответствующий код. Для этого откройте Запуск | Редактирование конфигураций и добавьте опцию
-Dkotlinx.coroutines.debugдля виртуальной машины:
Имя сопрограммы будет прикреплено к имени потока во время выполнения
main()с этой опцией. Вы также можете изменить шаблон для запуска всех файлов Kotlin и включить эту опцию по умолчанию.
Теперь весь код выполняется в одной сопрограмме, сопрограмме "загрузить участников", указанной выше, обозначенной как @coroutine#1. Ожидая результата, вы не должны повторно использовать поток для отправки других запросов, так как код написан последовательно. Новый запрос отправляется только после получения предыдущего результата.
Функции приостановки справедливо обрабатывают поток и не блокируют его для "ожидания". Однако это пока не вносит никакой конкуретности.
Конкурентность
Kotlin корутины потребляют значительно меньше ресурсов, чем потоки. Каждый раз, когда вы хотите запустить новое асинхронное вычисление, можно создать новую корутину.
Для запуска новой корутины используйте один из основных строителей корутин: launch, async, или runBlocking. Различные библиотеки могут определять дополнительные строители корутин.
async запускает новую корутину и возвращает объект Deferred. Deferred представляет концепцию, известную также как Future или Promise. Она хранит вычисление, но откладывает момент получения конечного результата; она обещает результат в будущем.
Основное различие между async и launch заключается в том, что launch используется для запуска вычисления, от которого не ожидается возврат конкретного результата. launch возвращает Job, представляющий корутину. Можно дождаться её завершения, вызвав Job.join()
Deferred — это универсальный тип, расширяющий Job. Вызов async может вернуть Deferred<Int> или Deferred<CustomType>, в зависимости от того, что возвращает лямбда (последнее выражение внутри лямбды — результат).
Чтобы получить результат корутины, можно вызвать await() на экземпляре Deferred. Во время ожидания результата корутина, из которой вызван этот await(), приостанавливается:
import kotlinx.coroutines.*
fun main() = runBlocking {
val deferred: Deferred<Int> = async {
loadData()
}
println("waiting...")
println(deferred.await())
}
suspend fun loadData(): Int {
println("loading...")
delay(1000L)
println("loaded!")
return 42
}
runBlocking используется в качестве моста между обычными и отложенными функциями, или между блокирующими и неблокирующими мирами. Она работает как адаптер для запуска главной корутины верхнего уровня. Она предназначена в основном для использования в функциях main() и тестах.
Если есть список объектов отложенных задач, можно вызвать awaitAll() для ожидания результатов всех из них:
import kotlinx.coroutines.*
fun main() = runBlocking {
val deferreds: List<Deferred<Int>> = (1..3).map {
async {
delay(1000L * it)
println("Loading $it")
it
}
}
val sum = deferreds.awaitAll().sum()
println("$sum")
}
Когда каждый запрос "contributors" запускается в новой корутине, все запросы запускаются асинхронно. Новый запрос можно отправить, прежде чем получить результат предыдущего:

Общее время загрузки примерно такое же, как в версии CALLBACKS, но без необходимости в обратных вызовах. Более того, async явно выделяет те части кода, которые выполняются конкуретнно.
Задача 5
В файле Request5Concurrent.kt реализуйте функцию loadContributorsConcurrent(), используя предыдущую функцию loadContributorsSuspend()
Подсказка к задаче 5
Вы можете запускать новые корутины только внутри области корутин. Скопируйте содержимое из loadContributorsSuspend() в вызов coroutineScope, чтобы вы могли вызывать там функции async
suspend fun loadContributorsConcurrent(
service: GitHubService,
req: RequestData
): List<User> = coroutineScope {
// ...
}
Основывайте свое решение на следующей схеме:
val deferreds: List<Deferred<List<User>>> = repos.map { repo ->
async {
// load contributors for each repo
}
}
deferreds.awaitAll() // List<List<User>>
Решение для задачи 5
Оборачивайте каждый запрос "contributors" с помощью async, чтобы создать столько корутин, сколько репозиториев. async возвращает Deferred<List<User>>. Это не проблема, так как создание новых корутин не сильно нагружает ресурсы, поэтому вы можете создавать столько, сколько нужно.
-
Вы больше не можете использовать
flatMap, так как результатmapтеперь список объектовDeferred, а не список списков.awaitAll()возвращаетList<List<User>>, поэтому вызовитеflatten().aggregate()для получения результата:suspend fun loadContributorsConcurrent( service: GitHubService, req: RequestData ): List<User> = coroutineScope { val repos = service .getOrgRepos(req.org) .also { logRepos(req, it) } .bodyList() val deferreds: List<Deferred<List<User>>> = repos.map { repo -> async { service.getRepoContributors(req.org, repo.name) .also { logUsers(repo, it) } .bodyList() } } deferreds.awaitAll().flatten().aggregate() } Запустите код и проверьте лог. Все корутины по-прежнему выполняются в главном потоке пользовательского интерфейса, так как многопоточность еще не использовалась, но вы уже можете увидеть преимущества одновременного выполнения корутин.
-
Чтобы изменить этот код на выполнение корутин "contributors" в разных потоках из общего пула потоков, укажите
Dispatchers.Defaultв качестве аргумента контекста для функцииasync:async(Dispatchers.Default) { }CoroutineDispatcherопределяет, в каком потоке или потоках должна выполняться соответствующая корутина. Если вы не укажете его как аргумент,asyncбудет использовать диспетчер из внешней области видимости.Dispatchers.Defaultпредставляет собой общий пул потоков на JVM. Этот пул обеспечивает средства для параллельного выполнения. Он состоит из столько потоков, сколько доступно ядер процессора, но при одном ядре все равно будет два потока.
-
Измените код в функции
loadContributorsConcurrent()для запуска новых корутин в разных потоках из общего пула потоков. Также добавьте дополнительное логирование перед отправкой запроса:async(Dispatchers.Default) { log("starting loading for ${repo.name}") service.getRepoContributors(req.org, repo.name) .also { logUsers(repo, it) } .bodyList() } -
Запустите программу еще раз. В журнале вы увидите, что каждая корутина может быть запущена в одном потоке из пула потоков и продолжена в другом:
1946 [DefaultDispatcher-worker-2 @coroutine#4] INFO Contributors - starting loading for kotlin-koans 1946 [DefaultDispatcher-worker-3 @coroutine#5] INFO Contributors - starting loading for dokka 1946 [DefaultDispatcher-worker-1 @coroutine#3] INFO Contributors - starting loading for ts2kt ... 2178 [DefaultDispatcher-worker-1 @coroutine#4] INFO Contributors - kotlin-koans: loaded 45 contributors 2569 [DefaultDispatcher-worker-1 @coroutine#5] INFO Contributors - dokka: loaded 36 contributors 2821 [DefaultDispatcher-worker-2 @coroutine#3] INFO Contributors - ts2kt: loaded 11 contributors
Например, в этом фрагменте лога
coroutine#4запущена в потокеworker-2и продолжена в потокеworker-1
В src/contributors/Contributors.kt, проверьте реализацию опции CONCURRENT:
-
Чтобы запустить корутину только в главном потоке пользовательского интерфейса, укажите
Dispatchers.Mainв качестве аргумента:launch(Dispatchers.Main) { updateResults() }Если главный поток занят, когда вы запускаете новую корутину в нём, корутина приостанавливается и планируется для выполнения в этом потоке. Корутина возобновится только когда поток освободится.
Хорошей практикой считается использование диспетчера из внешней области видимости, а не явное указание его в каждом конечной точке. Если вы определяете
loadContributorsConcurrent()без передачиDispatchers.Defaultв качестве аргумента, вы можете вызвать эту функцию в любом контексте: с диспетчеромDefault, с главным потоком пользовательского интерфейса или с пользовательским диспетчером.Как вы увидите позже, при вызове
loadContributorsConcurrent()из тестов, вы можете вызвать её в контексте сTestCoroutineDispatcher, что упрощает тестирование. Это делает решение более гибким.
-
Чтобы указать диспетчер со стороны вызывающей стороны, внесите следующие изменения в проект, позволяя
loadContributorsConcurrentзапускать корутины в унаследованном контексте:launch(Dispatchers.Default) { val users = loadContributorsConcurrent(service, req) withContext(Dispatchers.Main) { updateResults(users, startTime) } }updateResults()должна быть вызвана в главном потоке пользовательского интерфейса, поэтому вызывайте её с контекстомDispatchers.MainwithContext()вызывает заданный код с указанным контекстом корутины, приостанавливается до завершения и возвращает результат. Альтернативный, но более подробный способ выразить это — запустить новую корутину и явно дождаться (приостановив) её завершения:launch(context) { ... }.join()
Запустите код и убедитесь, что корутины выполняются в потоках из пула потоков.
Структурная конвейность
Область области корутины отвечает за структуру и отношения «родитель-потомок» между различными корутинами. Новые корутины обычно необходимо запускать внутри области.
Контекст корутины хранит дополнительную техническую информацию, используемую для запуска заданной корутины, например, пользовательское имя корутины или диспетчер, определяющий потоки, на которых корутина должна планироваться.
Когда launch, async, или runBlocking используются для запуска новой корутины, они автоматически создают соответствующую область. Все эти функции принимают лямбда-выражение с получателем в качестве аргумента, а CoroutineScope — это неявный тип получателя:
launch { /* this: CoroutineScope */ }
Новые корутины могут запускаться только внутри области.
launchиasyncобъявлены как расширения дляCoroutineScope, поэтому неявный или явный получатель всегда должен передаваться при их вызове.Исключением является корутина, запущенная с помощью
runBlocking, посколькуrunBlockingопределена как функция верхнего уровня. Но поскольку она блокирует текущий поток, она в основном предназначена для использования в функцияхmain()и тестах в качестве функции-моста.
Новая корутина внутри runBlocking, launch, или async запускается автоматически внутри области:
import kotlinx.coroutines.*
fun main() = runBlocking { /* this: CoroutineScope */
launch { /* ... */ }
// the same as:
this.launch { /* ... */ }
}
При вызове launch внутри runBlocking, она вызывается как расширение для неявного получателя типа CoroutineScope. В качестве альтернативы можно явно написать this.launch.
Вложенная корутина (запущенная с помощью launch в этом примере) может рассматриваться как потомок внешней корутины (запущенной с помощью runBlocking). Это отношение «родитель-потомок» реализуется через области; дочерняя корутина запускается из области, соответствующей родительской корутине.
Можно создать новую область без запуска новой корутины, используя функцию coroutineScope. Чтобы структурированно запустить новые корутины внутри функции suspend без доступа к внешней области, можно создать новую область корутины, которая автоматически становится дочерней областью внешней области, из которой вызывается функция suspend. loadContributorsConcurrent() — хороший пример.
Вы также можете запустить новую корутину из глобальной области, используя GlobalScope.async или GlobalScope.launch. Это создаст независимую корутину верхнего уровня.
Механизм, лежащий в основе структуры корутин, называется структурной конвейностью. Он предоставляет следующие преимущества по сравнению с глобальными областями:
Область, как правило, отвечает за дочерние корутины, чья продолжительность связана с продолжительностью области.
Область может автоматически отменять дочерние корутины, если что-то пойдет не так или пользователь передумает и решит аннулировать операцию.
Область автоматически ожидает завершения всех дочерних корутин. Поэтому, если область соответствует корутине, родительская корутина не завершается, пока не завершатся все корутины, запущенные в её области.
При использовании GlobalScope.async, нет структуры, которая связывает несколько корутин с более узкой областью. Корутины, запущенные из глобальной области, все независимы — их продолжительность ограничена только продолжительностью всего приложения. Можно сохранить ссылку на запущенную из глобальной области корутину и ожидать её завершения или отменить её явно, но это не произойдёт автоматически, как в случае со структурной конвейностью.
Отмена загрузки участников
Рассмотрим два варианта функции loadContributorsConcurrent(). Первый использует coroutineScope для запуска всех дочерних корутин, а второй — GlobalScope. Сравните, как ведут себя оба варианта при попытке отмены родительской корутины.
Скопируйте реализацию
loadContributorsConcurrent()изRequest5Concurrent.ktвloadContributorsNotCancellable()вRequest5NotCancellable.kt, а затем удалите создание новойcoroutineScope.-
Вызовы
asyncтеперь не разрешаются, поэтому запустите их с помощьюGlobalScope.async.suspend fun loadContributorsNotCancellable( service: GitHubService, req: RequestData ): List<User> { // #1 // ... GlobalScope.async { // #2 log("starting loading for ${repo.name}") // load repo contributors } // ... return deferreds.awaitAll().flatten().aggregate() // #3 }Функция теперь возвращает результат напрямую, а не в качестве последнего выражения внутри лямбда-выражения (строки
#1и#3).Все корутины «участников» запускаются внутри
GlobalScope, а не в качестве дочерних корутин области корутины (строка#2).
-
Добавьте задержку в 3 секунды во все корутины, отправляющие запросы, чтобы у вас было достаточно времени для отмены загрузки после запуска корутин, но до отправки запросов:
suspend fun loadContributorsConcurrent( service: GitHubService, req: RequestData ): List<User> = coroutineScope { // ... GlobalScope.async { log("starting loading for ${repo.name}") delay(3000) // load repo contributors } // ... } Запустите программу и выберите вариант КОНКУРЕНТНЫЙ для загрузки участников.
-
Подождите, пока все корутины «участников» не запустятся, а затем нажмите Отмена. В журнале нет новых результатов, что означает, что все запросы действительно были отменены:
2896 [AWT-EventQueue-0 @coroutine#1] INFO Contributors - kotlin: loaded 40 repos 2901 [DefaultDispatcher-worker-2 @coroutine#4] INFO Contributors - starting loading for kotlin-koans ... 2909 [DefaultDispatcher-worker-5 @coroutine#36] INFO Contributors - starting loading for mpp-example /* click on 'cancel' */ /* no requests are sent */
-
Повторите шаг 5, но на этот раз выберите вариант
NOT_CANCELLABLE.2570 [AWT-EventQueue-0 @coroutine#1] INFO Contributors - kotlin: loaded 30 repos 2579 [DefaultDispatcher-worker-1 @coroutine#4] INFO Contributors - starting loading for kotlin-koans ... 2586 [DefaultDispatcher-worker-6 @coroutine#36] INFO Contributors - starting loading for mpp-example /* click on 'cancel' */ /* but all the requests are still sent: */ 6402 [DefaultDispatcher-worker-5 @coroutine#4] INFO Contributors - kotlin-koans: loaded 45 contributors ... 9555 [DefaultDispatcher-worker-8 @coroutine#36] INFO Contributors - mpp-example: loaded 8 contributors
В этом случае корутины не отменяются, и все запросы по-прежнему отправляются.
-
Проверьте, как вызывается отмена в программе «участники». Когда нажимается кнопка Отмена, основная корутина «загрузки» явно отменяется, а дочерние корутины отменяются автоматически:
interface Contributors { fun loadContributors() { // ... when (getSelectedVariant()) { CONCURRENT -> { launch { val users = loadContributorsConcurrent(service, req) updateResults(users, startTime) }.setUpCancellation() // #1 } } } private fun Job.setUpCancellation() { val loadingJob = this // #2 // cancel the loading job if the 'cancel' button was clicked: val listener = ActionListener { loadingJob.cancel() // #3 updateLoadingStatus(CANCELED) } // add a listener to the 'cancel' button: addCancelListener(listener) // update the status and remove the listener // after the loading job is completed } }
Функция launch возвращает экземпляр Job. Job хранит ссылку на корутину «загрузки», которая загружает все данные и обновляет результаты. Вы можете вызвать функцию-расширение setUpCancellation() для неё (строка #1), передавая экземпляр Job как получатель.
Другой способ выразить это — явно написать:
val job = launch { }
job.setUpCancellation()
Для большей читабельности вы можете обратиться к получателю функции
setUpCancellation()внутри функции с новой переменнойloadingJob(строка#2).Затем вы можете добавить обработчик нажатия кнопки Отмена, чтобы при нажатии корутина
loadingJobотменялась (строка#3).
При использовании структурной конвейности достаточно отменить родительскую корутину, и это автоматически распространит отмену на все дочерние корутины.
Использование контекста внешней области
При запуске новых корутин в данной области гораздо проще убедиться, что все они выполняются с тем же контекстом. Также намного проще заменить контекст при необходимости.
Теперь давайте разберём, как работает использование диспетчера из внешней области. Новая область, созданная с помощью coroutineScope или с помощью конструкторов корутин, всегда наследует контекст из внешней области. В данном случае внешняя область — это область, из которой была вызвана функция suspend loadContributorsConcurrent().
launch(Dispatchers.Default) { // outer scope
val users = loadContributorsConcurrent(service, req)
// ...
}
Все вложенные корутины автоматически запускаются с унаследованным контекстом. Диспетчер является частью этого контекста. Вот почему все корутины, запущенные с помощью async, запускаются с контекстом диспетчера по умолчанию:
suspend fun loadContributorsConcurrent(
service: GitHubService, req: RequestData
): List<User> = coroutineScope {
// this scope inherits the context from the outer scope
// ...
async { // nested coroutine started with the inherited context
// ...
}
// ...
}
При использовании структурной конвейности вы можете указать основные элементы контекста (например, диспетчер) один раз при создании корутины верхнего уровня. Затем все вложенные корутины наследуют контекст и изменяют его только при необходимости.
Отображение прогресса
Несмотря на то, что информация для некоторых репозиториев загружается довольно быстро, пользователь видит результирующий список только после загрузки всех данных. До этого момента значок загрузчика отображает прогресс, но нет информации о текущем состоянии или о том, какие авторы уже загружены.
Вы можете отображать промежуточные результаты раньше и отображать всех авторов после загрузки данных для каждого репозитория:

Для реализации этой функциональности в src/tasks/Request6Progress.kt, вам нужно передать логику обновления пользовательского интерфейса в качестве обратного вызова, чтобы он вызывался на каждом промежуточном шаге:
suspend fun loadContributorsProgress(
service: GitHubService,
req: RequestData,
updateResults: suspend (List<User>, completed: Boolean) -> Unit
) {
// loading the data
// calling `updateResults()` on intermediate states
}
В месте вызова в Contributors.kt, обратный вызов передается для обновления результатов из потока Main для параметра PROGRESS:
launch(Dispatchers.Default) {
loadContributorsProgress(service, req) { users, completed ->
withContext(Dispatchers.Main) {
updateResults(users, startTime, completed)
}
}
}
Параметр
updateResults()объявлен какsuspendвloadContributorsProgress(). Необходимо вызватьwithContext, которая является функциейsuspendвнутри соответствующего лямбда-аргумента.Обратный вызов
updateResults()принимает дополнительный булевый параметр в качестве аргумента, указывающий, завершена ли загрузка и являются ли результаты окончательными.
Задача 6
В файле Request6Progress.kt реализуйте функцию loadContributorsProgress(), которая показывает промежуточный прогресс. Используйте функцию loadContributorsSuspend() из Request4Suspend.kt.
Используйте простую версию без конкурентного выполнения; вы добавите её позже в следующем разделе.
Промежуточный список авторов должен отображаться в «агрегированном» состоянии, а не только список пользователей, загруженных для каждого репозитория.
Общее количество вкладов каждого пользователя должно увеличиваться при загрузке данных для каждого нового репозитория.
Решение задачи 6
Для хранения промежуточного списка загруженных авторов в «агрегированном» состоянии определите переменную allUsers, которая хранит список пользователей, и затем обновляйте её после загрузки авторов для каждого нового репозитория:
suspend fun loadContributorsProgress(
service: GitHubService,
req: RequestData,
updateResults: suspend (List<User>, completed: Boolean) -> Unit
) {
val repos = service
.getOrgRepos(req.org)
.also { logRepos(req, it) }
.bodyList()
var allUsers = emptyList<User>()
for ((index, repo) in repos.withIndex()) {
val users = service.getRepoContributors(req.org, repo.name)
.also { logUsers(repo, it) }
.bodyList()
allUsers = (allUsers + users).aggregate()
updateResults(allUsers, index == repos.lastIndex)
}
}
Последовательное против одновременного выполнения
Обратный вызов updateResults() вызывается после завершения каждого запроса:

Этот код не включает конкурентного выполнения. Он последовательный, поэтому вам не нужна синхронизация.
Лучшим вариантом было бы отправлять запросы одновременно и обновлять промежуточные результаты после получения ответа для каждого репозитория:

Для добавления одновременного выполнения используйте каналы.
Каналы
Написание кода с общей изменяемой переменной довольно сложно и склонно к ошибкам (как в решении с обратными вызовами). Более простой способ — передавать информацию через коммуникацию, а не через общую изменяемую переменную. Корутины могут общаться друг с другом через каналы.
Каналы — это примитивы для обмена данными, позволяющие передавать данные между корутинами. Одна корутина может отправить некоторую информацию в канал, а другая может получить эту информацию из него:
Корутина, которая отправляет (производит) информацию, часто называется производителем, а корутина, которая получает (потребляет) информацию, называется потребителем. Одна или несколько корутин могут отправлять информацию в один и тот же канал, и одна или несколько корутин могут получать данные из него:
Когда много корутин получают информацию из одного и того же канала, каждый элемент обрабатывается только один раз одной из корутин-потребителей. После обработки элемента он немедленно удаляется из канала.
Вы можете представить канал как похожий на коллекцию элементов, или точнее, на очередь, в которую элементы добавляются с одного конца и извлекаются с другого. Однако есть важное различие: в отличие от коллекций, даже в их синхронизированных версиях, канал может приостановить send() и receive() операции. Это происходит, когда канал пуст или заполнен. Канал может быть заполнен, если размер канала имеет верхнюю границу.
Channel представлен тремя различными интерфейсами: SendChannel, ReceiveChannel, и Channel, причем последний расширяет два первых. Обычно вы создаете канал и передаете его производителям в качестве SendChannel экземпляра, чтобы только они могли отправлять информацию в канал. Вы передаете канал потребителям в качестве ReceiveChannel экземпляра, чтобы только они могли получать из него. Оба send и receive метода объявлены как suspend.
interface SendChannel<in E> {
suspend fun send(element: E)
fun close(): Boolean
}
interface ReceiveChannel<out E> {
suspend fun receive(): E
}
interface Channel<E> : SendChannel<E>, ReceiveChannel<E>
Производитель может закрыть канал, чтобы указать, что больше элементов не поступает.
В библиотеке определены несколько типов каналов. Они отличаются тем, сколько элементов они могут хранить внутри и приостанавливается ли вызов send() или нет. Для всех типов каналов вызов receive() работает аналогично: он получает элемент, если канал не пуст; в противном случае он приостанавливается.
- Канал без ограничений
-
Канал без ограничений — самый близкий аналог очереди: производители могут отправлять элементы в этот канал, и он будет расти неограниченно. Вызов
send()никогда не будет приостановлен. Если программа исчерпает память, вы получитеOutOfMemoryException. Разница между каналом без ограничений и очередью заключается в том, что когда потребитель пытается получить из пустого канала, он приостанавливается до тех пор, пока не будут отправлены новые элементы. - Буферизованный канал
-
Размер буферизованного канала ограничен указанным числом. Производители могут отправлять элементы в этот канал до тех пор, пока не будет достигнут предел размера. Все элементы хранятся внутри. Когда канал заполнен, следующий вызов `send` на нем приостанавливается, пока не освободится больше места.
- Канал согласования
-
Канал «согласования» — это канал без буфера, такой же, как буферизованный канал с нулевым размером. Одна из функций (
send()илиreceive()) всегда приостанавливается, пока не будет вызвана другая.Если вызов функции
send()и нет приостановленного вызоваreceive, готового обработать элемент, тоsend()приостанавливается. Аналогично, если вызов функцииreceiveи канал пуст, или, другими словами, нет приостановленного вызоваsend(), готового отправить элемент, вызовreceive()приостанавливается.Название «согласования» («встреча в согласованное время и месте») относится к тому факту, что
send()иreceive()должны «встретиться вовремя». - Слитый канал
-
Новый элемент, отправленный в слитый канал, перезапишет ранее отправленный элемент, поэтому получатель всегда получит только последний элемент. Вызов
send()никогда не приостанавливается.
При создании канала укажите его тип или размер буфера (если вам нужен буферизованный канал):
val rendezvousChannel = Channel<String>() val bufferedChannel = Channel<String>(10) val conflatedChannel = Channel<String>(CONFLATED) val unlimitedChannel = Channel<String>(UNLIMITED)
По умолчанию создается канал «согласования».
В следующем задании вы создадите канал «согласования», две корутины-производителя и одну корутину-потребителя:
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.*
fun main() = runBlocking<Unit> {
val channel = Channel<String>()
launch {
channel.send("A1")
channel.send("A2")
log("A done")
}
launch {
channel.send("B1")
log("B done")
}
launch {
repeat(3) {
val x = channel.receive()
log(x)
}
}
}
fun log(message: Any?) {
println("[${Thread.currentThread().name}] $message")
}
Задача 7
В src/tasks/Request7Channels.kt, реализуйте функцию loadContributorsChannels(), которая запрашивает всех участников GitHub одновременно и показывает промежуточный прогресс.
Используйте предыдущие функции, loadContributorsConcurrent() из Request5Concurrent.kt и loadContributorsProgress() из Request6Progress.kt.
Подсказка к задаче 7
Разные корутины, которые одновременно получают списки участников для разных репозиториев, могут отправлять все полученные результаты в один и тот же канал:
val channel = Channel<List<User>>()
for (repo in repos) {
launch {
val users = TODO()
// ...
channel.send(users)
}
}
Затем элементы из этого канала могут быть получены по одному и обработаны:
repeat(repos.size) {
val users = channel.receive()
// ...
}
Поскольку вызовы receive() последовательны, дополнительная синхронизация не требуется.
Решение задачи 7
Как и в функции loadContributorsProgress(), вы можете создать переменную allUsers для хранения промежуточных состояний списка «все участники». Каждый новый список, полученный из канала, добавляется в список всех пользователей. Результат агрегируется и состояние обновляется с помощью обратного вызова updateResults:
suspend fun loadContributorsChannels(
service: GitHubService,
req: RequestData,
updateResults: suspend (List<User>, completed: Boolean) -> Unit
) = coroutineScope {
val repos = service
.getOrgRepos(req.org)
.also { logRepos(req, it) }
.bodyList()
val channel = Channel<List<User>>()
for (repo in repos) {
launch {
val users = service.getRepoContributors(req.org, repo.name)
.also { logUsers(repo, it) }
.bodyList()
channel.send(users)
}
}
var allUsers = emptyList<User>()
repeat(repos.size) {
val users = channel.receive()
allUsers = (allUsers + users).aggregate()
updateResults(allUsers, it == repos.lastIndex)
}
}
Результаты для разных репозиториев добавляются в канал как только они готовы. Сначала, когда все запросы отправлены, и данные не получены, вызов
receive()приостанавливается. В этом случае вся корутина «загрузить участников» приостанавливается.Затем, когда список пользователей отправлен в канал, корутина «загрузить участников» возобновляется, вызов
receive()возвращает этот список, и результаты немедленно обновляются.
Теперь вы можете запустить программу и выбрать опцию КАНАЛЫ для загрузки участников и просмотра результата.
Хотя ни корутины, ни каналы не устраняют полностью сложность, связанную с конкурентностью, они облегчают жизнь, когда нужно понять, что происходит.
Тестирование сопрограмм
Теперь давайте протестируем все решения, чтобы убедиться, что решение с конкурентными сопрограммами быстрее, чем решение с функциями suspend, и что решение с каналами быстрее, чем простое решение «прогресса».
В следующей задаче вы сравните общее время выполнения решений. Вы смоделируете службу GitHub и заставите эту службу возвращать результаты после заданных таймаутов:
repos request - returns an answer within 1000 ms delay repo-1 - 1000 ms delay repo-2 - 1200 ms delay repo-3 - 800 ms delay
Последовательное решение с функциями suspend должно занимать около 4000 мс (4000 = 1000 + (1000 + 1200 + 800)). Конкурентное решение должно занимать около 2200 мс (2200 = 1000 + max(1000, 1200, 800)).
Для решений, которые отображают прогресс, вы также можете проверить промежуточные результаты с отметками времени.
Соответствующие тестовые данные определены в test/contributors/testData.kt, а файлы Request4SuspendKtTest, Request7ChannelsKtTest, и т. д. содержат простые тесты, которые используют вызовы моделируемой службы.
Однако здесь есть две проблемы:
Эти тесты занимают слишком много времени для выполнения. Каждый тест занимает около 2-4 секунд, и вам нужно ждать результатов каждый раз. Это не очень эффективно.
Вы не можете полагаться на точное время выполнения решения, потому что оно по-прежнему занимает дополнительное время для подготовки и запуска кода. Вы можете добавить константу, но тогда время будет отличаться от машины к машине. Задержки моделируемой службы должны быть больше этой константы, чтобы вы могли увидеть разницу. Если константа составляет 0,5 секунды, то задержки в 0,1 секунды будут недостаточны.
Лучше использовать специальные фреймворки для тестирования времени, выполняя тот же код несколько раз (что увеличивает общее время еще больше), но это сложно освоить и настроить.
Чтобы решить эти проблемы и убедиться, что решения с предоставленными задержками тестов ведут себя ожидаемым образом, одно быстрее другого, используйте виртуальное время с особым диспетчером тестов. Этот диспетчер отслеживает прошедшее виртуальное время с начала и выполняет все сразу в реальном времени. Когда вы запускаете сопрограммы на этом диспетчере, delay вернется сразу и продвинет виртуальное время.
Тесты, использующие этот механизм, выполняются быстро, но вы все равно можете проверить, что происходит в разные моменты виртуального времени. Общее время выполнения резко уменьшается:

Для использования виртуального времени замените вызов runBlocking на runBlockingTest. runBlockingTest принимает расширяющую лямбду для TestCoroutineScope в качестве аргумента. Когда вы вызываете delay в функции suspend внутри этого специального контекста, delay увеличит виртуальное время вместо задержки в реальном времени:
@Test
fun testDelayInSuspend() = runBlockingTest {
val realStartTime = System.currentTimeMillis()
val virtualStartTime = currentTime
foo()
println("${System.currentTimeMillis() - realStartTime} ms") // ~ 6 ms
println("${currentTime - virtualStartTime} ms") // 1000 ms
}
suspend fun foo() {
delay(1000) // auto-advances without delay
println("foo") // executes eagerly when foo() is called
}
Текущее виртуальное время можно проверить, используя свойство currentTime объекта TestCoroutineScope.
Фактическое время выполнения в этом примере составляет несколько миллисекунд, в то время как виртуальное время равно аргументу задержки, который равен 1000 миллисекундам.
Чтобы получить полный эффект «виртуального» delay в дочерних сопрограммах, запустите все дочерние сопрограммы с помощью TestCoroutineDispatcher. В противном случае это не сработает. Этот диспетчер автоматически наследуется от других TestCoroutineScope, если вы не укажете другой диспетчер:
@Test
fun testDelayInLaunch() = runBlockingTest {
val realStartTime = System.currentTimeMillis()
val virtualStartTime = currentTime
bar()
println("${System.currentTimeMillis() - realStartTime} ms") // ~ 11 ms
println("${currentTime - virtualStartTime} ms") // 1000 ms
}
suspend fun bar() = coroutineScope {
launch {
delay(1000) // auto-advances without delay
println("bar") // executes eagerly when bar() is called
}
}
Если launch вызывается в контексте Dispatchers.Default в приведенном выше примере, тест завершится сбоем. Вы получите исключение, указывающее, что задача еще не завершена.
Вы можете протестировать функцию loadContributorsConcurrent() таким образом только в том случае, если она запускает дочерние сопрограммы с унаследованным контекстом без его изменения с помощью диспетчера Dispatchers.Default.
Вы можете указывать элементы контекста, такие как диспетчер, при вызове функции, а не при ее определении, что обеспечивает большую гибкость и простоту тестирования.
По умолчанию компилятор выводит предупреждения, если вы используете экспериментальный API тестирования. Чтобы подавить эти предупреждения, добавьте аннотацию к функции теста или всему классу, содержащему тесты, с помощью @OptIn(ExperimentalCoroutinesApi::class). Добавьте аргумент компилятора, указывающий компилятору, что вы используете экспериментальный API:
compileTestKotlin {
kotlinOptions {
freeCompilerArgs += "-Xuse-experimental=kotlin.Experimental"
}
}
В проекте, соответствующем этому учебнику, аргумент компилятора уже добавлен в скрипт Gradle.
Задача 8
Переработайте следующие тесты в tests/tasks/ с использованием виртуального времени вместо реального:
Request4SuspendKtTest.kt
Request5ConcurrentKtTest.kt
Request6ProgressKtTest.kt
Request7ChannelsKtTest.kt
Сравните общее время выполнения до и после применения переработки.
Подсказка к задаче 8
-
Замените вызов
runBlockingнаrunBlockingTest, аSystem.currentTimeMillis()наcurrentTime.@Test fun test() = runBlockingTest { val startTime = currentTime // action val totalTime = currentTime - startTime // testing result } Разкомментируйте утверждения, проверяющие точное виртуальное время.
Не забудьте добавить
@UseExperimental(ExperimentalCoroutinesApi::class).
Решение задачи 8
Вот решения для случаев с конкурентностью и каналами:
fun testConcurrent() = runBlockingTest {
val startTime = currentTime
val result = loadContributorsConcurrent(MockGithubService, testRequestData)
Assert.assertEquals("Wrong result for 'loadContributorsConcurrent'", expectedConcurrentResults.users, result)
val totalTime = currentTime - startTime
Assert.assertEquals(
"The calls run concurrently, so the total virtual time should be 2200 ms: " +
"1000 for repos request plus max(1000, 1200, 800) = 1200 for concurrent contributors requests)",
expectedConcurrentResults.timeFromStart, totalTime
)
}
Сначала проверьте, что результаты доступны ровно в ожидаемое виртуальное время, а затем проверьте сами результаты:
fun testChannels() = runBlockingTest {
val startTime = currentTime
var index = 0
loadContributorsChannels(MockGithubService, testRequestData) { users, _ ->
val expected = concurrentProgressResults[index++]
val time = currentTime - startTime
Assert.assertEquals(
"Expected intermediate results after ${expected.timeFromStart} ms:",
expected.timeFromStart, time
)
Assert.assertEquals("Wrong intermediate results after $time:", expected.users, users)
}
}
Первый промежуточный результат для последней версии с каналами становится доступным раньше, чем для версии с прогрессом, и вы можете увидеть разницу в тестах, использующих виртуальное время.
Что дальше
Посетите мастер-класс по асинхронному программированию с Kotlin на KotlinConf.
Узнайте больше о использовании виртуального времени и экспериментального пакета тестирования.
© 2010–2022 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/coroutines-and-channels.html